From c930fb45f1e6736b83bdbefad0bf8c2b4ce33c61 Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 10 Sep 2024 17:23:47 +0800 Subject: [PATCH] test --- ecstore/src/store.rs | 33 ++++++-- rustfs/src/main.rs | 12 ++- rustfs/src/storage/ecfs.rs | 159 +++++++++++++++++++++++++++++-------- 3 files changed, 162 insertions(+), 42 deletions(-) diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 487886d4e..5116d5624 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -1,7 +1,6 @@ use crate::{ bucket_meta::BucketMetadata, disk::{error::DiskError, DeleteOptions, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET}, - disks_layout::DisksLayout, endpoints::EndpointServerPools, error::{Error, Result}, peer::{PeerS3Client, S3PeerSys}, @@ -16,11 +15,31 @@ use crate::{ use futures::future::join_all; use http::HeaderMap; use s3s::{dto::StreamingBlob, Body}; -use std::collections::{HashMap, HashSet}; +use std::{ + collections::{HashMap, HashSet}, + sync::Arc, +}; use time::OffsetDateTime; +use tokio::sync::Mutex; use tracing::{debug, warn}; use uuid::Uuid; +use lazy_static::lazy_static; + +lazy_static! { + pub static ref GLOBAL_OBJECT_API: Arc>> = Arc::new(Mutex::new(None)); +} + +pub fn new_object_layer_fn() -> Arc>> { + // 这里不需要显式地锁定和解锁,因为 Arc 提供了必要的线程安全性 + GLOBAL_OBJECT_API.clone() +} + +async fn set_object_layer(o: ECStore) { + let mut global_object_api = GLOBAL_OBJECT_API.lock().await; + *global_object_api = Some(o); +} + #[derive(Debug)] pub struct ECStore { pub id: uuid::Uuid, @@ -32,7 +51,7 @@ pub struct ECStore { } impl ECStore { - pub async fn new(address: String, endpoint_pools: EndpointServerPools) -> Result { + pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<()> { // let layouts = DisksLayout::try_from(endpoints.as_slice())?; let mut deployment_id = None; @@ -97,13 +116,17 @@ impl ECStore { let peer_sys = S3PeerSys::new(&endpoint_pools, local_disks.clone()); - Ok(ECStore { + let ec = ECStore { id: deployment_id.unwrap(), disk_map, pools, local_disks, peer_sys, - }) + }; + + set_object_layer(ec).await; + + Ok(()) } fn single_pool(&self) -> bool { diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index 3a8fbba69..94cca1cf6 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -2,7 +2,7 @@ mod config; mod storage; use clap::Parser; -use ecstore::{disks_layout::DisksLayout, endpoints::EndpointServerPools, error::Result}; +use ecstore::{endpoints::EndpointServerPools, error::Result, store::ECStore}; use hyper_util::{ rt::{TokioExecutor, TokioIo}, server::conn::auto::Builder as ConnBuilder, @@ -10,7 +10,7 @@ use hyper_util::{ use s3s::{auth::SimpleAuth, service::S3ServiceBuilder}; use std::{io::IsTerminal, net::SocketAddr, str::FromStr}; use tokio::net::TcpListener; -use tracing::{debug, info}; +use tracing::{debug, info, warn}; use tracing_error::ErrorLayer; use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt}; @@ -66,12 +66,14 @@ async fn run(opt: config::Opt) -> Result<()> { // }) // }; + // 用于rpc let (endpoint_pools, _) = EndpointServerPools::from_volumes(opt.address.clone().as_str(), opt.volumes.clone())?; // Setup S3 service // 本项目使用s3s库来实现s3服务 let service = { - let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new(opt.address.clone(), endpoint_pools).await?); + // let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new(opt.address.clone(), endpoint_pools).await?); + let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new().await?); //设置AK和SK //其中部份内容从config配置文件中读取 let mut access_key = String::from_str(config::DEFAULT_ACCESS_KEY).unwrap(); @@ -135,6 +137,10 @@ async fn run(opt: config::Opt) -> Result<()> { }); } + warn!(" init store"); + // init store + ECStore::new(opt.address.clone(), endpoint_pools.clone()).await?; + tokio::select! { () = graceful.shutdown() => { tracing::debug!("Gracefully shutdown!"); diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 9de4d419a..f7a50bae3 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -1,5 +1,6 @@ use bytes::Bytes; use ecstore::disk::error::DiskError; +use ecstore::store::new_object_layer_fn; use ecstore::store_api::BucketOptions; use ecstore::store_api::CompletePart; use ecstore::store_api::HTTPRangeSpec; @@ -24,9 +25,7 @@ use std::str::FromStr; use transform_stream::AsyncTryStream; use uuid::Uuid; -use ecstore::endpoints::EndpointServerPools; use ecstore::error::Result; -use ecstore::store::ECStore; use tracing::debug; macro_rules! try_ { @@ -42,13 +41,13 @@ macro_rules! try_ { #[derive(Debug)] pub struct FS { - pub store: ECStore, + // pub store: ECStore, } impl FS { - pub async fn new(address: String, endpoint_pools: EndpointServerPools) -> Result { - let store: ECStore = ECStore::new(address, endpoint_pools).await?; - Ok(Self { store }) + pub async fn new() -> Result { + // let store: ECStore = ECStore::new(address, endpoint_pools).await?; + Ok(Self {}) } } #[async_trait::async_trait] @@ -61,8 +60,15 @@ impl S3 for FS { async fn create_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!( - self.store + store .make_bucket(&input.bucket, &MakeBucketOptions { force_create: true }) .await ); @@ -87,7 +93,13 @@ impl S3 for FS { async fn delete_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - try_!(self.store.delete_bucket(&input.bucket).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!(store.delete_bucket(&input.bucket).await); Ok(S3Response::new(DeleteBucketOutput {})) } @@ -112,7 +124,13 @@ impl S3 for FS { let objects: Vec = vec![dobj]; - let (dobjs, _errs) = try_!(self.store.delete_objects(&bucket, objects, ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await); // TODO: let errors; @@ -173,7 +191,14 @@ impl S3 for FS { }) .collect(); - let (dobjs, _errs) = try_!(self.store.delete_objects(&bucket, objects, ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await); // info!("delete_objects res {:?} {:?}", &dobjs, errs); let deleted = dobjs @@ -207,7 +232,14 @@ impl S3 for FS { // mc get 1 let input = req.input; - if let Err(e) = self.store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -243,11 +275,14 @@ impl S3 for FS { let h = HeaderMap::new(); let opts = &ObjectOptions::default(); - let reader = try_!( - self.store - .get_object_reader(bucket.as_str(), key.as_str(), range, h, opts) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let reader = try_!(store.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts).await); let info = reader.object_info; @@ -270,7 +305,14 @@ impl S3 for FS { async fn head_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - if let Err(e) = self.store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -287,7 +329,14 @@ impl S3 for FS { // mc get 2 let HeadObjectInput { bucket, key, .. } = req.input; - let info = try_!(self.store.get_object_info(&bucket, &key, &ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let info = try_!(store.get_object_info(&bucket, &key, &ObjectOptions::default()).await); debug!("info {:?}", info); let content_type = try_!(ContentType::from_str("application/x-msdownload")); @@ -307,7 +356,14 @@ impl S3 for FS { async fn list_buckets(&self, _: S3Request) -> S3Result> { // mc ls - let bucket_infos = try_!(self.store.list_bucket(&BucketOptions {}).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let bucket_infos = try_!(store.list_bucket(&BucketOptions {}).await); let buckets: Vec = bucket_infos .iter() @@ -355,8 +411,15 @@ impl S3 for FS { let prefix = prefix.unwrap_or_default(); let delimiter = delimiter.unwrap_or_default(); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let object_infos = try_!( - self.store + store .list_objects_v2( &bucket, &prefix, @@ -435,9 +498,16 @@ impl S3 for FS { let reader = PutObjReader::new(body, content_length as usize); - try_!(self.store.put_object(&bucket, &key, reader, &ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; - // self.store.put_object(bucket, object, data, opts); + try_!(store.put_object(&bucket, &key, reader, &ObjectOptions::default()).await); + + // store.put_object(bucket, object, data, opts); let output = PutObjectOutput { ..Default::default() }; Ok(S3Response::new(output)) @@ -456,11 +526,15 @@ impl S3 for FS { debug!("create_multipart_upload meta {:?}", &metadata); - let MultipartUploadResult { upload_id, .. } = try_!( - self.store - .new_multipart_upload(&bucket, &key, &ObjectOptions::default()) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let MultipartUploadResult { upload_id, .. } = + try_!(store.new_multipart_upload(&bucket, &key, &ObjectOptions::default()).await); let output = CreateMultipartUploadOutput { bucket: Some(bucket), @@ -495,11 +569,14 @@ impl S3 for FS { let data = PutObjReader::new(body, content_length as usize); let opts = ObjectOptions::default(); - try_!( - self.store - .put_object_part(&bucket, &key, &upload_id, part_id, data, &opts) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + try_!(store.put_object_part(&bucket, &key, &upload_id, part_id, data, &opts).await); let output = UploadPartOutput { ..Default::default() }; Ok(S3Response::new(output)) @@ -554,8 +631,15 @@ impl S3 for FS { uploaded_parts.push(CompletePart::from(part)); } + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!( - self.store + store .complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts) .await ); @@ -577,9 +661,16 @@ impl S3 for FS { bucket, key, upload_id, .. } = req.input; + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let opts = &ObjectOptions::default(); try_!( - self.store + store .abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts) .await );