This commit is contained in:
weisd
2024-09-10 17:23:47 +08:00
parent f7e7f6e257
commit c930fb45f1
3 changed files with 162 additions and 42 deletions
+28 -5
View File
@@ -1,7 +1,6 @@
use crate::{ use crate::{
bucket_meta::BucketMetadata, bucket_meta::BucketMetadata,
disk::{error::DiskError, DeleteOptions, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET}, disk::{error::DiskError, DeleteOptions, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET},
disks_layout::DisksLayout,
endpoints::EndpointServerPools, endpoints::EndpointServerPools,
error::{Error, Result}, error::{Error, Result},
peer::{PeerS3Client, S3PeerSys}, peer::{PeerS3Client, S3PeerSys},
@@ -16,11 +15,31 @@ use crate::{
use futures::future::join_all; use futures::future::join_all;
use http::HeaderMap; use http::HeaderMap;
use s3s::{dto::StreamingBlob, Body}; use s3s::{dto::StreamingBlob, Body};
use std::collections::{HashMap, HashSet}; use std::{
collections::{HashMap, HashSet},
sync::Arc,
};
use time::OffsetDateTime; use time::OffsetDateTime;
use tokio::sync::Mutex;
use tracing::{debug, warn}; use tracing::{debug, warn};
use uuid::Uuid; use uuid::Uuid;
use lazy_static::lazy_static;
lazy_static! {
pub static ref GLOBAL_OBJECT_API: Arc<Mutex<Option<ECStore>>> = Arc::new(Mutex::new(None));
}
pub fn new_object_layer_fn() -> Arc<Mutex<Option<ECStore>>> {
// 这里不需要显式地锁定和解锁,因为 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)] #[derive(Debug)]
pub struct ECStore { pub struct ECStore {
pub id: uuid::Uuid, pub id: uuid::Uuid,
@@ -32,7 +51,7 @@ pub struct ECStore {
} }
impl ECStore { impl ECStore {
pub async fn new(address: String, endpoint_pools: EndpointServerPools) -> Result<Self> { pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<()> {
// let layouts = DisksLayout::try_from(endpoints.as_slice())?; // let layouts = DisksLayout::try_from(endpoints.as_slice())?;
let mut deployment_id = None; let mut deployment_id = None;
@@ -97,13 +116,17 @@ impl ECStore {
let peer_sys = S3PeerSys::new(&endpoint_pools, local_disks.clone()); let peer_sys = S3PeerSys::new(&endpoint_pools, local_disks.clone());
Ok(ECStore { let ec = ECStore {
id: deployment_id.unwrap(), id: deployment_id.unwrap(),
disk_map, disk_map,
pools, pools,
local_disks, local_disks,
peer_sys, peer_sys,
}) };
set_object_layer(ec).await;
Ok(())
} }
fn single_pool(&self) -> bool { fn single_pool(&self) -> bool {
+9 -3
View File
@@ -2,7 +2,7 @@ mod config;
mod storage; mod storage;
use clap::Parser; use clap::Parser;
use ecstore::{disks_layout::DisksLayout, endpoints::EndpointServerPools, error::Result}; use ecstore::{endpoints::EndpointServerPools, error::Result, store::ECStore};
use hyper_util::{ use hyper_util::{
rt::{TokioExecutor, TokioIo}, rt::{TokioExecutor, TokioIo},
server::conn::auto::Builder as ConnBuilder, server::conn::auto::Builder as ConnBuilder,
@@ -10,7 +10,7 @@ use hyper_util::{
use s3s::{auth::SimpleAuth, service::S3ServiceBuilder}; use s3s::{auth::SimpleAuth, service::S3ServiceBuilder};
use std::{io::IsTerminal, net::SocketAddr, str::FromStr}; use std::{io::IsTerminal, net::SocketAddr, str::FromStr};
use tokio::net::TcpListener; use tokio::net::TcpListener;
use tracing::{debug, info}; use tracing::{debug, info, warn};
use tracing_error::ErrorLayer; use tracing_error::ErrorLayer;
use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt}; 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())?; let (endpoint_pools, _) = EndpointServerPools::from_volumes(opt.address.clone().as_str(), opt.volumes.clone())?;
// Setup S3 service // Setup S3 service
// 本项目使用s3s库来实现s3服务 // 本项目使用s3s库来实现s3服务
let service = { 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 //设置AK和SK
//其中部份内容从config配置文件中读取 //其中部份内容从config配置文件中读取
let mut access_key = String::from_str(config::DEFAULT_ACCESS_KEY).unwrap(); 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! { tokio::select! {
() = graceful.shutdown() => { () = graceful.shutdown() => {
tracing::debug!("Gracefully shutdown!"); tracing::debug!("Gracefully shutdown!");
+125 -34
View File
@@ -1,5 +1,6 @@
use bytes::Bytes; use bytes::Bytes;
use ecstore::disk::error::DiskError; use ecstore::disk::error::DiskError;
use ecstore::store::new_object_layer_fn;
use ecstore::store_api::BucketOptions; use ecstore::store_api::BucketOptions;
use ecstore::store_api::CompletePart; use ecstore::store_api::CompletePart;
use ecstore::store_api::HTTPRangeSpec; use ecstore::store_api::HTTPRangeSpec;
@@ -24,9 +25,7 @@ use std::str::FromStr;
use transform_stream::AsyncTryStream; use transform_stream::AsyncTryStream;
use uuid::Uuid; use uuid::Uuid;
use ecstore::endpoints::EndpointServerPools;
use ecstore::error::Result; use ecstore::error::Result;
use ecstore::store::ECStore;
use tracing::debug; use tracing::debug;
macro_rules! try_ { macro_rules! try_ {
@@ -42,13 +41,13 @@ macro_rules! try_ {
#[derive(Debug)] #[derive(Debug)]
pub struct FS { pub struct FS {
pub store: ECStore, // pub store: ECStore,
} }
impl FS { impl FS {
pub async fn new(address: String, endpoint_pools: EndpointServerPools) -> Result<Self> { pub async fn new() -> Result<Self> {
let store: ECStore = ECStore::new(address, endpoint_pools).await?; // let store: ECStore = ECStore::new(address, endpoint_pools).await?;
Ok(Self { store }) Ok(Self {})
} }
} }
#[async_trait::async_trait] #[async_trait::async_trait]
@@ -61,8 +60,15 @@ impl S3 for FS {
async fn create_bucket(&self, req: S3Request<CreateBucketInput>) -> S3Result<S3Response<CreateBucketOutput>> { async fn create_bucket(&self, req: S3Request<CreateBucketInput>) -> S3Result<S3Response<CreateBucketOutput>> {
let input = req.input; 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_!( try_!(
self.store store
.make_bucket(&input.bucket, &MakeBucketOptions { force_create: true }) .make_bucket(&input.bucket, &MakeBucketOptions { force_create: true })
.await .await
); );
@@ -87,7 +93,13 @@ impl S3 for FS {
async fn delete_bucket(&self, req: S3Request<DeleteBucketInput>) -> S3Result<S3Response<DeleteBucketOutput>> { async fn delete_bucket(&self, req: S3Request<DeleteBucketInput>) -> S3Result<S3Response<DeleteBucketOutput>> {
let input = req.input; 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 {})) Ok(S3Response::new(DeleteBucketOutput {}))
} }
@@ -112,7 +124,13 @@ impl S3 for FS {
let objects: Vec<ObjectToDelete> = vec![dobj]; let objects: Vec<ObjectToDelete> = 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; // TODO: let errors;
@@ -173,7 +191,14 @@ impl S3 for FS {
}) })
.collect(); .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); // info!("delete_objects res {:?} {:?}", &dobjs, errs);
let deleted = dobjs let deleted = dobjs
@@ -207,7 +232,14 @@ impl S3 for FS {
// mc get 1 // mc get 1
let input = req.input; 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) { if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket)); return Err(s3_error!(NoSuchBucket));
} else { } else {
@@ -243,11 +275,14 @@ impl S3 for FS {
let h = HeaderMap::new(); let h = HeaderMap::new();
let opts = &ObjectOptions::default(); let opts = &ObjectOptions::default();
let reader = try_!( let layer = new_object_layer_fn();
self.store let lock = layer.lock().await;
.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts) let store = match lock.as_ref() {
.await 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; let info = reader.object_info;
@@ -270,7 +305,14 @@ impl S3 for FS {
async fn head_bucket(&self, req: S3Request<HeadBucketInput>) -> S3Result<S3Response<HeadBucketOutput>> { async fn head_bucket(&self, req: S3Request<HeadBucketInput>) -> S3Result<S3Response<HeadBucketOutput>> {
let input = req.input; 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) { if DiskError::VolumeNotFound.is(&e) {
return Err(s3_error!(NoSuchBucket)); return Err(s3_error!(NoSuchBucket));
} else { } else {
@@ -287,7 +329,14 @@ impl S3 for FS {
// mc get 2 // mc get 2
let HeadObjectInput { bucket, key, .. } = req.input; 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); debug!("info {:?}", info);
let content_type = try_!(ContentType::from_str("application/x-msdownload")); let content_type = try_!(ContentType::from_str("application/x-msdownload"));
@@ -307,7 +356,14 @@ impl S3 for FS {
async fn list_buckets(&self, _: S3Request<ListBucketsInput>) -> S3Result<S3Response<ListBucketsOutput>> { async fn list_buckets(&self, _: S3Request<ListBucketsInput>) -> S3Result<S3Response<ListBucketsOutput>> {
// mc ls // 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> = bucket_infos let buckets: Vec<Bucket> = bucket_infos
.iter() .iter()
@@ -355,8 +411,15 @@ impl S3 for FS {
let prefix = prefix.unwrap_or_default(); let prefix = prefix.unwrap_or_default();
let delimiter = delimiter.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_!( let object_infos = try_!(
self.store store
.list_objects_v2( .list_objects_v2(
&bucket, &bucket,
&prefix, &prefix,
@@ -435,9 +498,16 @@ impl S3 for FS {
let reader = PutObjReader::new(body, content_length as usize); 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() }; let output = PutObjectOutput { ..Default::default() };
Ok(S3Response::new(output)) Ok(S3Response::new(output))
@@ -456,11 +526,15 @@ impl S3 for FS {
debug!("create_multipart_upload meta {:?}", &metadata); debug!("create_multipart_upload meta {:?}", &metadata);
let MultipartUploadResult { upload_id, .. } = try_!( let layer = new_object_layer_fn();
self.store let lock = layer.lock().await;
.new_multipart_upload(&bucket, &key, &ObjectOptions::default()) let store = match lock.as_ref() {
.await 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 { let output = CreateMultipartUploadOutput {
bucket: Some(bucket), bucket: Some(bucket),
@@ -495,11 +569,14 @@ impl S3 for FS {
let data = PutObjReader::new(body, content_length as usize); let data = PutObjReader::new(body, content_length as usize);
let opts = ObjectOptions::default(); let opts = ObjectOptions::default();
try_!( let layer = new_object_layer_fn();
self.store let lock = layer.lock().await;
.put_object_part(&bucket, &key, &upload_id, part_id, data, &opts) let store = match lock.as_ref() {
.await 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() }; let output = UploadPartOutput { ..Default::default() };
Ok(S3Response::new(output)) Ok(S3Response::new(output))
@@ -554,8 +631,15 @@ impl S3 for FS {
uploaded_parts.push(CompletePart::from(part)); 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_!( try_!(
self.store store
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts) .complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts)
.await .await
); );
@@ -577,9 +661,16 @@ impl S3 for FS {
bucket, key, upload_id, .. bucket, key, upload_id, ..
} = req.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",))),
};
let opts = &ObjectOptions::default(); let opts = &ObjectOptions::default();
try_!( try_!(
self.store store
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts) .abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
.await .await
); );