mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-19 11:06:17 +00:00
refactor: centralize external ECStore facade aliases (#3746)
This commit is contained in:
@@ -16,7 +16,7 @@
|
||||
#![allow(dead_code)]
|
||||
|
||||
use async_trait::async_trait;
|
||||
use rustfs_ecstore::api::rpc::{self as ecstore_rpc, TonicInterceptor};
|
||||
use rustfs_ecstore::api::rpc as ecstore_rpc;
|
||||
use rustfs_lock::{
|
||||
LockClient, LockError, LockId, LockInfo, LockRequest, LockResponse, LockStats, LockStatus, LockType, Result,
|
||||
types::{LockMetadata, LockPriority},
|
||||
@@ -25,6 +25,8 @@ use rustfs_protos::proto_gen::node_service::{BatchGenerallyLockRequest, Generall
|
||||
use tonic::Request;
|
||||
use tracing::{info, warn};
|
||||
|
||||
type TonicInterceptor = ecstore_rpc::TonicInterceptor;
|
||||
|
||||
/// gRPC lock client without authentication for testing
|
||||
/// Similar to RemoteClient but uses no_auth client
|
||||
#[derive(Debug, Clone)]
|
||||
|
||||
@@ -16,10 +16,8 @@
|
||||
use crate::common::workspace_root;
|
||||
use futures::future::join_all;
|
||||
use rmp_serde::{Deserializer, Serializer};
|
||||
use rustfs_ecstore::api::{
|
||||
disk::{VolumeInfo, WalkDirOptions},
|
||||
rpc::{self as ecstore_rpc, TonicInterceptor},
|
||||
};
|
||||
use rustfs_ecstore::api::disk as ecstore_disk;
|
||||
use rustfs_ecstore::api::rpc as ecstore_rpc;
|
||||
use rustfs_filemeta::{MetaCacheEntry, MetacacheReader, MetacacheWriter};
|
||||
use rustfs_protos::proto_gen::node_service::WalkDirRequest;
|
||||
use rustfs_protos::{
|
||||
@@ -38,6 +36,10 @@ use tonic::codegen::tokio_stream::StreamExt;
|
||||
|
||||
const CLUSTER_ADDR: &str = "http://localhost:9000";
|
||||
|
||||
type TonicInterceptor = ecstore_rpc::TonicInterceptor;
|
||||
type VolumeInfo = ecstore_disk::VolumeInfo;
|
||||
type WalkDirOptions = ecstore_disk::WalkDirOptions;
|
||||
|
||||
fn signature_interceptor() -> TonicInterceptor {
|
||||
TonicInterceptor::Signature(ecstore_rpc::gen_tonic_signature_interceptor())
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
|
||||
use aws_sdk_s3::{Client, Config};
|
||||
use http::header::{CONTENT_TYPE, HOST};
|
||||
use reqwest::StatusCode;
|
||||
use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys;
|
||||
use rustfs_ecstore::api::bucket as ecstore_bucket;
|
||||
use rustfs_madmin::{
|
||||
AddServiceAccountReq, ListServiceAccountsResp, PeerInfo, PeerSite, ReplicateAddStatus, ReplicateEditStatus,
|
||||
ReplicateRemoveStatus, SRRemoveReq, SRResyncOpStatus, SRStatusInfo, SiteReplicationInfo, SyncStatus,
|
||||
@@ -37,6 +37,7 @@ use time::Duration as TimeDuration;
|
||||
use tokio::time::{Duration, sleep};
|
||||
|
||||
type TestResult = Result<(), Box<dyn Error + Send + Sync>>;
|
||||
type BucketTargetSys = ecstore_bucket::bucket_target_sys::BucketTargetSys;
|
||||
|
||||
#[derive(Debug, Clone, serde::Deserialize)]
|
||||
struct ReplicationResetStatusResponse {
|
||||
|
||||
@@ -14,15 +14,19 @@
|
||||
|
||||
//! test endpoint index settings
|
||||
|
||||
use rustfs_ecstore::api::{
|
||||
disk::endpoint::Endpoint,
|
||||
layout::{EndpointServerPools, Endpoints, PoolEndpoints},
|
||||
storage::{ECStore, init_local_disks},
|
||||
};
|
||||
use rustfs_ecstore::api::disk as ecstore_disk;
|
||||
use rustfs_ecstore::api::layout as ecstore_layout;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use std::net::SocketAddr;
|
||||
use tempfile::TempDir;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
type ECStore = ecstore_storage::ECStore;
|
||||
type Endpoint = ecstore_disk::endpoint::Endpoint;
|
||||
type EndpointServerPools = ecstore_layout::EndpointServerPools;
|
||||
type Endpoints = ecstore_layout::Endpoints;
|
||||
type PoolEndpoints = ecstore_layout::PoolEndpoints;
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn test_endpoint_index_settings() -> anyhow::Result<()> {
|
||||
let temp_dir = TempDir::new()?;
|
||||
@@ -74,7 +78,7 @@ async fn test_endpoint_index_settings() -> anyhow::Result<()> {
|
||||
}
|
||||
|
||||
// test ECStore initialization
|
||||
init_local_disks(endpoint_pools.clone()).await?;
|
||||
ecstore_storage::init_local_disks(endpoint_pools.clone()).await?;
|
||||
|
||||
let server_addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
|
||||
let ecstore = ECStore::new(server_addr, endpoint_pools, CancellationToken::new()).await?;
|
||||
|
||||
@@ -12,13 +12,16 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use rustfs_ecstore::api::disk::{DiskStore, endpoint::Endpoint};
|
||||
use rustfs_ecstore::api::disk as ecstore_disk;
|
||||
use rustfs_heal::heal::{
|
||||
event::{HealEvent, Severity},
|
||||
task::{HealPriority, HealType},
|
||||
utils,
|
||||
};
|
||||
|
||||
type DiskStore = ecstore_disk::DiskStore;
|
||||
type Endpoint = ecstore_disk::endpoint::Endpoint;
|
||||
|
||||
#[test]
|
||||
fn test_heal_event_to_heal_request_no_panic() {
|
||||
// Test that invalid pool/set indices don't cause panic
|
||||
|
||||
@@ -14,12 +14,10 @@
|
||||
|
||||
use http::HeaderMap;
|
||||
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
|
||||
use rustfs_ecstore::api::{
|
||||
bucket::metadata_sys::init_bucket_metadata_sys,
|
||||
disk::endpoint::Endpoint,
|
||||
layout::{EndpointServerPools, Endpoints, PoolEndpoints},
|
||||
storage::{ECStore, init_local_disks},
|
||||
};
|
||||
use rustfs_ecstore::api::bucket as ecstore_bucket;
|
||||
use rustfs_ecstore::api::disk as ecstore_disk;
|
||||
use rustfs_ecstore::api::layout as ecstore_layout;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use rustfs_heal::heal::{
|
||||
manager::{HealConfig, HealManager},
|
||||
storage::{ECStoreHealStorage, HealObjectOptions as ObjectOptions, HealPutObjReader as PutObjReader, HealStorageAPI},
|
||||
@@ -41,6 +39,12 @@ const HEAL_FORMAT_WAIT_TIMEOUT: Duration = Duration::from_secs(25);
|
||||
const HEAL_FORMAT_WAIT_INTERVAL: Duration = Duration::from_millis(250);
|
||||
const NON_INLINE_TEST_DATA_SIZE: usize = 256 * 1024 + 137;
|
||||
|
||||
type ECStore = ecstore_storage::ECStore;
|
||||
type Endpoint = ecstore_disk::endpoint::Endpoint;
|
||||
type EndpointServerPools = ecstore_layout::EndpointServerPools;
|
||||
type Endpoints = ecstore_layout::Endpoints;
|
||||
type PoolEndpoints = ecstore_layout::PoolEndpoints;
|
||||
|
||||
fn non_inline_test_data() -> Vec<u8> {
|
||||
(0..NON_INLINE_TEST_DATA_SIZE).map(|idx| (idx % 251) as u8).collect()
|
||||
}
|
||||
@@ -123,7 +127,7 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>, Arc<ECStoreHealStorage
|
||||
let endpoint_pools = EndpointServerPools::from(vec![pool_endpoints]);
|
||||
|
||||
// format disks (only first time)
|
||||
init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
ecstore_storage::init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
|
||||
// create ECStore with dynamic port 0 (let OS assign) or fixed 9001 if free
|
||||
let port = 9001; // for simplicity
|
||||
@@ -141,7 +145,7 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>, Arc<ECStoreHealStorage
|
||||
.await
|
||||
.unwrap();
|
||||
let buckets = buckets_list.into_iter().map(|v| v.name).collect();
|
||||
init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
ecstore_bucket::metadata_sys::init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
|
||||
// Create heal storage layer
|
||||
let heal_storage = Arc::new(ECStoreHealStorage::new(ecstore.clone()));
|
||||
|
||||
@@ -15,10 +15,11 @@
|
||||
use crate::error::{Error, Result};
|
||||
use manager::IamCache;
|
||||
use oidc::OidcSys;
|
||||
use rustfs_ecstore::api::{
|
||||
config as ecstore_config, error as ecstore_error, global as ecstore_global, notification as ecstore_notification,
|
||||
storage as ecstore_storage,
|
||||
};
|
||||
use rustfs_ecstore::api::config as ecstore_config;
|
||||
use rustfs_ecstore::api::error as ecstore_error;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::notification as ecstore_notification;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use std::sync::{Arc, OnceLock};
|
||||
use store::object::ObjectStore;
|
||||
use sys::IamSys;
|
||||
|
||||
@@ -20,7 +20,9 @@ use rustfs_config::notify::{
|
||||
NOTIFY_POSTGRES_SUB_SYS, NOTIFY_PULSAR_SUB_SYS, NOTIFY_REDIS_SUB_SYS, NOTIFY_WEBHOOK_SUB_SYS,
|
||||
};
|
||||
use rustfs_config::server_config::{Config, KVS};
|
||||
use rustfs_ecstore::api::{config as ecstore_config, global as ecstore_global, storage as ecstore_storage};
|
||||
use rustfs_ecstore::api::config as ecstore_config;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use rustfs_targets::{Target, arn::TargetID};
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
@@ -29,10 +29,12 @@ use crate::metrics::collectors::{
|
||||
use chrono::Utc;
|
||||
use rustfs_common::heal_channel::HealScanMode;
|
||||
use rustfs_common::metrics::global_metrics;
|
||||
use rustfs_ecstore::api::{
|
||||
bucket as ecstore_bucket, capacity as ecstore_capacity, data_usage as ecstore_data_usage, error as ecstore_error,
|
||||
global as ecstore_global, storage as ecstore_storage,
|
||||
};
|
||||
use rustfs_ecstore::api::bucket as ecstore_bucket;
|
||||
use rustfs_ecstore::api::capacity as ecstore_capacity;
|
||||
use rustfs_ecstore::api::data_usage as ecstore_data_usage;
|
||||
use rustfs_ecstore::api::error as ecstore_error;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use rustfs_iam::{get_global_iam_sys, oidc::oidc_plugin_authn_metrics_snapshot};
|
||||
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
|
||||
use rustfs_io_metrics::{ProcessStatusSnapshot, snapshot_process_resource_and_system};
|
||||
|
||||
@@ -35,9 +35,10 @@
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use rustfs_ecstore::api::{
|
||||
bucket as ecstore_bucket, error as ecstore_error, global as ecstore_global, storage as ecstore_storage,
|
||||
};
|
||||
use rustfs_ecstore::api::bucket as ecstore_bucket;
|
||||
use rustfs_ecstore::api::error as ecstore_error;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
|
||||
pub mod account;
|
||||
pub mod acl;
|
||||
|
||||
@@ -25,9 +25,10 @@ use object_store::{
|
||||
};
|
||||
use pin_project_lite::pin_project;
|
||||
use rustfs_common::DEFAULT_DELIMITER;
|
||||
use rustfs_ecstore::api::{
|
||||
error as ecstore_error, global as ecstore_global, set_disk as ecstore_set_disk, storage as ecstore_storage,
|
||||
};
|
||||
use rustfs_ecstore::api::error as ecstore_error;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::set_disk as ecstore_set_disk;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use rustfs_storage_api::{HTTPRangeSpec, ObjectIO as _, ObjectOperations as _};
|
||||
use s3s::S3Result;
|
||||
use s3s::dto::SelectObjectContentInput;
|
||||
|
||||
@@ -12,29 +12,20 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use ecstore_bucket::metadata_sys;
|
||||
use ecstore_disk::DiskAPI as _;
|
||||
use ecstore_global::GLOBAL_TierConfigMgr;
|
||||
use ecstore_tier::warm_backend::{WarmBackend as ScannerWarmBackend, build_transition_put_options};
|
||||
use futures::FutureExt;
|
||||
use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT;
|
||||
use rustfs_ecstore::api::{
|
||||
bucket::{
|
||||
lifecycle::{
|
||||
bucket_lifecycle_ops::{enqueue_transition_for_existing_objects, init_background_expiry},
|
||||
lifecycle::TransitionOptions,
|
||||
},
|
||||
metadata::BUCKET_LIFECYCLE_CONFIG,
|
||||
metadata_sys::{self, init_bucket_metadata_sys},
|
||||
versioning_sys::BucketVersioningSys,
|
||||
},
|
||||
capacity::path2_bucket_object_with_base_path,
|
||||
client::transition_api::{ReadCloser, ReaderImpl},
|
||||
disk::{DiskAPI as _, DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, new_disk},
|
||||
global::GLOBAL_TierConfigMgr,
|
||||
layout::{EndpointServerPools, Endpoints, PoolEndpoints},
|
||||
storage::{ECStore, init_local_disks},
|
||||
tier::{
|
||||
tier_config::{TierConfig, TierMinIO, TierType},
|
||||
warm_backend::{WarmBackend as ScannerWarmBackend, WarmBackendGetOpts, build_transition_put_options},
|
||||
},
|
||||
};
|
||||
use rustfs_ecstore::api::bucket as ecstore_bucket;
|
||||
use rustfs_ecstore::api::capacity as ecstore_capacity;
|
||||
use rustfs_ecstore::api::client as ecstore_client;
|
||||
use rustfs_ecstore::api::disk as ecstore_disk;
|
||||
use rustfs_ecstore::api::global as ecstore_global;
|
||||
use rustfs_ecstore::api::layout as ecstore_layout;
|
||||
use rustfs_ecstore::api::storage as ecstore_storage;
|
||||
use rustfs_ecstore::api::tier as ecstore_tier;
|
||||
use rustfs_filemeta::FileMeta;
|
||||
use rustfs_scanner::scanner_folder::ScannerItem;
|
||||
use rustfs_scanner::scanner_io::ScannerIODisk;
|
||||
@@ -66,6 +57,23 @@ use uuid::Uuid;
|
||||
static GLOBAL_ENV: OnceLock<(Vec<PathBuf>, Arc<ECStore>)> = OnceLock::new();
|
||||
static INIT: Once = Once::new();
|
||||
const TRANSITION_WAIT_TIMEOUT: Duration = Duration::from_secs(15);
|
||||
const BUCKET_LIFECYCLE_CONFIG: &str = ecstore_bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
|
||||
const STORAGE_FORMAT_FILE: &str = ecstore_disk::STORAGE_FORMAT_FILE;
|
||||
|
||||
type BucketVersioningSys = ecstore_bucket::versioning_sys::BucketVersioningSys;
|
||||
type ECStore = ecstore_storage::ECStore;
|
||||
type Endpoint = ecstore_disk::endpoint::Endpoint;
|
||||
type EndpointServerPools = ecstore_layout::EndpointServerPools;
|
||||
type Endpoints = ecstore_layout::Endpoints;
|
||||
type PoolEndpoints = ecstore_layout::PoolEndpoints;
|
||||
type DiskOption = ecstore_disk::DiskOption;
|
||||
type ReadCloser = ecstore_client::transition_api::ReadCloser;
|
||||
type ReaderImpl = ecstore_client::transition_api::ReaderImpl;
|
||||
type TierConfig = ecstore_tier::tier_config::TierConfig;
|
||||
type TierMinIO = ecstore_tier::tier_config::TierMinIO;
|
||||
type TierType = ecstore_tier::tier_config::TierType;
|
||||
type TransitionOptions = ecstore_bucket::lifecycle::lifecycle::TransitionOptions;
|
||||
type WarmBackendGetOpts = ecstore_tier::warm_backend::WarmBackendGetOpts;
|
||||
|
||||
fn init_tracing() {
|
||||
INIT.call_once(|| {
|
||||
@@ -125,7 +133,7 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>) {
|
||||
let endpoint_pools = EndpointServerPools::from(vec![pool_endpoints]);
|
||||
|
||||
// format disks (only first time)
|
||||
init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
ecstore_storage::init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
|
||||
// create ECStore with dynamic port 0 (let OS assign) or fixed 9002 if free
|
||||
let port = 9002; // for simplicity
|
||||
@@ -143,10 +151,10 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>) {
|
||||
.await
|
||||
.unwrap();
|
||||
let buckets = buckets_list.into_iter().map(|v| v.name).collect();
|
||||
init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
ecstore_bucket::metadata_sys::init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
|
||||
// Initialize background expiry workers
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
|
||||
// Store in global once lock
|
||||
let _ = GLOBAL_ENV.set((disk_paths.clone(), ecstore.clone()));
|
||||
@@ -194,7 +202,7 @@ async fn setup_isolated_test_env(init_expiry: bool) -> (Vec<PathBuf>, Arc<ECStor
|
||||
};
|
||||
|
||||
let endpoint_pools = EndpointServerPools::from(vec![pool_endpoints]);
|
||||
init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
ecstore_storage::init_local_disks(endpoint_pools.clone()).await.unwrap();
|
||||
|
||||
let server_addr: std::net::SocketAddr = "127.0.0.1:0".parse().unwrap();
|
||||
let ecstore = ECStore::new(server_addr, endpoint_pools, CancellationToken::new())
|
||||
@@ -209,10 +217,10 @@ async fn setup_isolated_test_env(init_expiry: bool) -> (Vec<PathBuf>, Arc<ECStor
|
||||
.await
|
||||
.unwrap();
|
||||
let buckets = buckets_list.into_iter().map(|v| v.name).collect();
|
||||
init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
ecstore_bucket::metadata_sys::init_bucket_metadata_sys(ecstore.clone(), buckets).await;
|
||||
|
||||
if init_expiry {
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
}
|
||||
|
||||
(disk_paths, ecstore)
|
||||
@@ -507,7 +515,7 @@ async fn free_version_count(disk_path: &Path, bucket: &str, object: &str) -> usi
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(0);
|
||||
let disk = new_disk(
|
||||
let disk = ecstore_disk::new_disk(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
cleanup: false,
|
||||
@@ -591,7 +599,7 @@ async fn scan_object_with_lifecycle(disk_path: &Path, bucket: &str, object: &str
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(0);
|
||||
let disk = new_disk(
|
||||
let disk = ecstore_disk::new_disk(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
cleanup: false,
|
||||
@@ -602,7 +610,8 @@ async fn scan_object_with_lifecycle(disk_path: &Path, bucket: &str, object: &str
|
||||
.expect("failed to open local disk");
|
||||
let metadata_path = disk_path.join(bucket).join(object).join(STORAGE_FORMAT_FILE);
|
||||
let relative_path = metadata_path.to_string_lossy().to_string();
|
||||
let (_, scanner_path) = path2_bucket_object_with_base_path(disk_path.to_string_lossy().as_ref(), relative_path.as_str());
|
||||
let (_, scanner_path) =
|
||||
ecstore_capacity::path2_bucket_object_with_base_path(disk_path.to_string_lossy().as_ref(), relative_path.as_str());
|
||||
let file_type = fs::metadata(&metadata_path)
|
||||
.await
|
||||
.expect("failed to stat object metadata")
|
||||
@@ -633,7 +642,7 @@ async fn scan_object_metadata(disk_path: &Path, bucket: &str, object: &str) {
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(0);
|
||||
let disk = new_disk(
|
||||
let disk = ecstore_disk::new_disk(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
cleanup: false,
|
||||
@@ -644,7 +653,8 @@ async fn scan_object_metadata(disk_path: &Path, bucket: &str, object: &str) {
|
||||
.expect("failed to open local disk");
|
||||
let metadata_path = disk_path.join(bucket).join(object).join(STORAGE_FORMAT_FILE);
|
||||
let relative_path = metadata_path.to_string_lossy().to_string();
|
||||
let (_, scanner_path) = path2_bucket_object_with_base_path(disk_path.to_string_lossy().as_ref(), relative_path.as_str());
|
||||
let (_, scanner_path) =
|
||||
ecstore_capacity::path2_bucket_object_with_base_path(disk_path.to_string_lossy().as_ref(), relative_path.as_str());
|
||||
let file_type = fs::metadata(&metadata_path)
|
||||
.await
|
||||
.expect("failed to stat object metadata")
|
||||
@@ -961,9 +971,12 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to upload transition metadata test object");
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), put_bucket.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned put object");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
put_bucket.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned put object");
|
||||
|
||||
let put_info = wait_for_transition(&ecstore, put_bucket.as_str(), put_object, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1031,9 +1044,12 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to complete multipart upload");
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), multipart_bucket.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned multipart object");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
multipart_bucket.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned multipart object");
|
||||
|
||||
let multipart_info = wait_for_transition(&ecstore, multipart_bucket.as_str(), multipart_object, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1082,9 +1098,12 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to copy object");
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), dst_bucket.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned copied object");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
dst_bucket.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned copied object");
|
||||
|
||||
let copy_info = wait_for_transition(&ecstore, dst_bucket.as_str(), dst_object, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1105,9 +1124,12 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to set lifecycle configuration");
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transition for existing objects");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
bucket_name.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transition for existing objects");
|
||||
|
||||
let info = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1182,9 +1204,12 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to complete multipart upload");
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned restore object");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
bucket_name.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned restore object");
|
||||
|
||||
let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1254,9 +1279,12 @@ mod serial_tests {
|
||||
.expect("Failed to set lifecycle configuration");
|
||||
|
||||
upload_test_object(&ecstore, bucket_name.as_str(), object_name, initial_payload).await;
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str())
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned object");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
bucket_name.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("Failed to enqueue transitioned object");
|
||||
|
||||
let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1278,7 +1306,7 @@ mod serial_tests {
|
||||
"stale transitioned remote object should still exist before scanner fallback runs"
|
||||
);
|
||||
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await;
|
||||
|
||||
assert!(
|
||||
@@ -1339,7 +1367,7 @@ mod serial_tests {
|
||||
"stale transitioned remote object should still exist before scanner cleanup runs"
|
||||
);
|
||||
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await;
|
||||
|
||||
assert!(
|
||||
@@ -1382,9 +1410,12 @@ mod serial_tests {
|
||||
let remote_object = transitioned.transitioned_object.name.clone();
|
||||
assert!(backend.objects.lock().await.contains_key(&remote_object));
|
||||
|
||||
enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str())
|
||||
.await
|
||||
.expect("existing-object backfill should succeed after compensation transition");
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects(
|
||||
ecstore.clone(),
|
||||
bucket_name.as_str(),
|
||||
)
|
||||
.await
|
||||
.expect("existing-object backfill should succeed after compensation transition");
|
||||
|
||||
let info = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT)
|
||||
.await
|
||||
@@ -1706,7 +1737,7 @@ mod serial_tests {
|
||||
|
||||
assert!(object_exists(&ecstore, bucket_name.as_str(), object_name).await);
|
||||
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await;
|
||||
|
||||
assert!(
|
||||
@@ -1810,7 +1841,7 @@ mod serial_tests {
|
||||
.await
|
||||
.expect("Failed to set noncurrent lifecycle configuration");
|
||||
|
||||
init_background_expiry(ecstore.clone()).await;
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await;
|
||||
|
||||
scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user