merge main

This commit is contained in:
weisd
2024-09-27 20:12:54 +08:00
43 changed files with 3798 additions and 737 deletions
+1
View File
@@ -9,6 +9,7 @@ rust-version.workspace = true
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[dependencies]
log.workspace = true
async-trait.workspace = true
bytes.workspace = true
clap.workspace = true
+168 -16
View File
@@ -1,33 +1,35 @@
use std::{error::Error, io::ErrorKind, pin::Pin};
use ecstore::{
disk::{DeleteOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, WalkDirOptions},
disk::{DeleteOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, UpdateMetadataOpts, WalkDirOptions},
erasure::{ReadAt, Write},
peer::{LocalPeerS3Client, PeerS3Client},
store::{all_local_disk_path, find_local_disk},
store_api::{BucketOptions, FileInfo, MakeBucketOptions},
store_api::{BucketOptions, DeleteBucketOptions, FileInfo, MakeBucketOptions},
};
use futures::{Stream, StreamExt};
use lock::{lock_args::LockArgs, Locker, GLOBAL_LOCAL_SERVER};
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tonic::{Request, Response, Status, Streaming};
use tracing::{debug, error, info};
use protos::{
models::{PingBody, PingBodyBuilder},
proto_gen::node_service::{
node_service_server::NodeService as Node, DeleteBucketRequest, DeleteBucketResponse, DeleteRequest, DeleteResponse,
DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DeleteVolumeResponse, GenerallyLockRequest,
GenerallyLockResponse, GetBucketInfoRequest, GetBucketInfoResponse, ListBucketRequest, ListBucketResponse,
ListDirRequest, ListDirResponse, ListVolumesRequest, ListVolumesResponse, MakeBucketRequest, MakeBucketResponse,
MakeVolumeRequest, MakeVolumeResponse, MakeVolumesRequest, MakeVolumesResponse, PingRequest, PingResponse,
ReadAllRequest, ReadAllResponse, ReadAtRequest, ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse,
ReadVersionRequest, ReadVersionResponse, ReadXlRequest, ReadXlResponse, RenameDataRequest, RenameDataResponse,
RenameFileRequst, RenameFileResponse, StatVolumeRequest, StatVolumeResponse, WalkDirRequest, WalkDirResponse,
WriteAllRequest, WriteAllResponse, WriteMetadataRequest, WriteMetadataResponse, WriteRequest, WriteResponse,
node_service_server::NodeService as Node, DeleteBucketRequest, DeleteBucketResponse, DeletePathsRequest,
DeletePathsResponse, DeleteRequest, DeleteResponse, DeleteVersionRequest, DeleteVersionResponse, DeleteVersionsRequest,
DeleteVersionsResponse, DeleteVolumeRequest, DeleteVolumeResponse, GenerallyLockRequest, GenerallyLockResponse,
GetBucketInfoRequest, GetBucketInfoResponse, ListBucketRequest, ListBucketResponse, ListDirRequest, ListDirResponse,
ListVolumesRequest, ListVolumesResponse, MakeBucketRequest, MakeBucketResponse, MakeVolumeRequest, MakeVolumeResponse,
MakeVolumesRequest, MakeVolumesResponse, PingRequest, PingResponse, ReadAllRequest, ReadAllResponse, ReadAtRequest,
ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse, ReadVersionRequest, ReadVersionResponse, ReadXlRequest,
ReadXlResponse, RenameDataRequest, RenameDataResponse, RenameFileRequst, RenameFileResponse, RenamePartRequst,
RenamePartResponse, StatVolumeRequest, StatVolumeResponse, UpdateMetadataRequest, UpdateMetadataResponse, WalkDirRequest,
WalkDirResponse, WriteAllRequest, WriteAllResponse, WriteMetadataRequest, WriteMetadataResponse, WriteRequest,
WriteResponse,
},
};
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tonic::{Request, Response, Status, Streaming};
use tracing::{debug, error, info};
type ResponseStream<T> = Pin<Box<dyn Stream<Item = Result<T, tonic::Status>> + Send>>;
@@ -208,7 +210,11 @@ impl Node for NodeService {
debug!("make bucket");
let request = request.into_inner();
match self.local_peer.delete_bucket(&request.bucket).await {
match self
.local_peer
.delete_bucket(&request.bucket, &DeleteBucketOptions { force: false })
.await
{
Ok(_) => Ok(tonic::Response::new(DeleteBucketResponse {
success: true,
error_info: None,
@@ -297,6 +303,36 @@ impl Node for NodeService {
}
}
async fn rename_part(&self, request: Request<RenamePartRequst>) -> Result<Response<RenamePartResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
match disk
.rename_part(
&request.src_volume,
&request.src_path,
&request.dst_volume,
&request.dst_path,
request.meta,
)
.await
{
Ok(_) => Ok(tonic::Response::new(RenamePartResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(RenamePartResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(RenamePartResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn rename_file(&self, request: Request<RenameFileRequst>) -> Result<Response<RenameFileResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
@@ -749,6 +785,68 @@ impl Node for NodeService {
}
}
async fn delete_paths(&self, request: Request<DeletePathsRequest>) -> Result<Response<DeletePathsResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let paths = request.paths.iter().map(|s| s.as_str()).collect::<Vec<&str>>();
match disk.delete_paths(&request.volume, &paths).await {
Ok(_) => Ok(tonic::Response::new(DeletePathsResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(DeletePathsResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(DeletePathsResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn update_metadata(&self, request: Request<UpdateMetadataRequest>) -> Result<Response<UpdateMetadataResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let file_info = match serde_json::from_str::<FileInfo>(&request.file_info) {
Ok(file_info) => file_info,
Err(_) => {
return Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not decode FileInfoVersions".to_string()),
}));
}
};
let opts = match serde_json::from_str::<UpdateMetadataOpts>(&request.opts) {
Ok(opts) => opts,
Err(_) => {
return Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not decode UpdateMetadataOpts".to_string()),
}));
}
};
match disk.update_metadata(&request.volume, &request.path, file_info, opts).await {
Ok(_) => Ok(tonic::Response::new(UpdateMetadataResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn write_metadata(&self, request: Request<WriteMetadataRequest>) -> Result<Response<WriteMetadataResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
@@ -854,6 +952,60 @@ impl Node for NodeService {
}
}
async fn delete_version(&self, request: Request<DeleteVersionRequest>) -> Result<Response<DeleteVersionResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let file_info = match serde_json::from_str::<FileInfo>(&request.file_info) {
Ok(file_info) => file_info,
Err(_) => {
return Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not decode FileInfoVersions".to_string()),
}));
}
};
let opts = match serde_json::from_str::<DeleteOptions>(&request.opts) {
Ok(opts) => opts,
Err(_) => {
return Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not decode DeleteOptions".to_string()),
}));
}
};
match disk
.delete_version(&request.volume, &request.path, file_info, request.force_del_marker, opts)
.await
{
Ok(raw_file_info) => match serde_json::to_string(&raw_file_info) {
Ok(raw_file_info) => Ok(tonic::Response::new(DeleteVersionResponse {
success: true,
raw_file_info,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some(err.to_string()),
})),
},
Err(err) => Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn delete_versions(&self, request: Request<DeleteVersionsRequest>) -> Result<Response<DeleteVersionsResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
+2 -3
View File
@@ -21,7 +21,7 @@ use service::hybrid;
use std::{io::IsTerminal, net::SocketAddr, str::FromStr};
use tokio::net::TcpListener;
use tonic::{metadata::MetadataValue, Request, Status};
use tracing::{debug, info, warn};
use tracing::{debug, info};
use tracing_error::ErrorLayer;
use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt};
@@ -178,12 +178,11 @@ async fn run(opt: config::Opt) -> Result<()> {
}
});
warn!(" init store");
// init store
ECStore::new(opt.address.clone(), endpoint_pools.clone())
.await
.map_err(|err| Error::from_string(err.to_string()))?;
warn!(" init store success!");
info!(" init store success!");
tokio::select! {
_ = tokio::signal::ctrl_c() => {
+268 -2
View File
@@ -1,8 +1,12 @@
use bytes::BufMut;
use bytes::Bytes;
use ecstore::bucket_meta::BucketMetadata;
use ecstore::disk::error::DiskError;
use ecstore::disk::RUSTFS_META_BUCKET;
use ecstore::store::new_object_layer_fn;
use ecstore::store_api::BucketOptions;
use ecstore::store_api::CompletePart;
use ecstore::store_api::DeleteBucketOptions;
use ecstore::store_api::HTTPRangeSpec;
use ecstore::store_api::MakeBucketOptions;
use ecstore::store_api::MultipartUploadResult;
@@ -15,6 +19,7 @@ use futures::{Stream, StreamExt};
use http::HeaderMap;
use s3s::dto::*;
use s3s::s3_error;
use s3s::Body;
use s3s::S3Error;
use s3s::S3ErrorCode;
use s3s::S3Result;
@@ -92,14 +97,18 @@ impl S3 for FS {
#[tracing::instrument(level = "debug", skip(self, req))]
async fn delete_bucket(&self, req: S3Request<DeleteBucketInput>) -> S3Result<S3Response<DeleteBucketOutput>> {
let input = req.input;
// TODO: DeleteBucketInput 没有force参数?
let layer = new_object_layer_fn();
let lock = layer.read().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);
try_!(
store
.delete_bucket(&input.bucket, &DeleteBucketOptions { force: false })
.await
);
Ok(S3Response::new(DeleteBucketOutput {}))
}
@@ -251,6 +260,7 @@ impl S3 for FS {
Ok(S3Response::new(output))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn get_object_lock_configuration(
&self,
_req: S3Request<GetObjectLockConfigurationInput>,
@@ -619,6 +629,7 @@ impl S3 for FS {
..
} = req.input;
// error!("complete_multipart_upload {:?}", multipart_upload);
// mc cp step 5
let Some(multipart_upload) = multipart_upload else { return Err(s3_error!(InvalidPart)) };
@@ -676,6 +687,261 @@ impl S3 for FS {
);
Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() }))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_bucket_tagging(&self, req: S3Request<PutBucketTaggingInput>) -> S3Result<S3Response<PutBucketTaggingOutput>> {
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
log::debug!("bucket: {bucket}, tagging: {tagging:?}");
// check bucket exists.
let _bucket = self
.head_bucket(S3Request::new(HeadBucketInput {
bucket: bucket.clone(),
expected_bucket_owner: None,
}))
.await?;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let meta_obj = try_!(
store
.get_object_reader(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
HTTPRangeSpec::nil(),
Default::default(),
&ObjectOptions::default(),
)
.await
);
let stream = meta_obj.stream;
let mut data = vec![];
pin_mut!(stream);
while let Some(x) = stream.next().await {
let x = try_!(x);
data.put_slice(&x[..]);
}
let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..]));
if tagging.tag_set.is_empty() {
meta.tagging = None;
} else {
meta.tagging = Some(tagging.tag_set.into_iter().map(|x| (x.key, x.value)).collect())
}
let data = try_!(meta.marshal_msg());
let len = data.len();
try_!(
store
.put_object(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
PutObjReader::new(StreamingBlob::from(Body::from(data)), len),
&ObjectOptions::default(),
)
.await
);
Ok(S3Response::new(Default::default()))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn get_bucket_tagging(&self, req: S3Request<GetBucketTaggingInput>) -> S3Result<S3Response<GetBucketTaggingOutput>> {
let GetBucketTaggingInput { bucket, .. } = req.input;
// check bucket exists.
let _bucket = self
.head_bucket(S3Request::new(HeadBucketInput {
bucket: bucket.clone(),
expected_bucket_owner: None,
}))
.await?;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let meta_obj = try_!(
store
.get_object_reader(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
HTTPRangeSpec::nil(),
Default::default(),
&ObjectOptions::default(),
)
.await
);
let stream = meta_obj.stream;
let mut data = vec![];
pin_mut!(stream);
while let Some(x) = stream.next().await {
let x = try_!(x);
data.put_slice(&x[..]);
}
let meta = try_!(BucketMetadata::unmarshal_from(&data[..]));
if meta.tagging.is_none() {
return Err({
let mut err = S3Error::with_message(S3ErrorCode::Custom("NoSuchTagSet".into()), "The TagSet does not exist");
err.set_status_code("404".try_into().unwrap());
err
});
}
Ok(S3Response::new(GetBucketTaggingOutput {
tag_set: meta
.tagging
.unwrap()
.into_iter()
.map(|(key, value)| Tag { key, value })
.collect(),
}))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn delete_bucket_tagging(
&self,
req: S3Request<DeleteBucketTaggingInput>,
) -> S3Result<S3Response<DeleteBucketTaggingOutput>> {
let DeleteBucketTaggingInput { bucket, .. } = req.input;
// check bucket exists.
let _bucket = self
.head_bucket(S3Request::new(HeadBucketInput {
bucket: bucket.clone(),
expected_bucket_owner: None,
}))
.await?;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let meta_obj = try_!(
store
.get_object_reader(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
HTTPRangeSpec::nil(),
Default::default(),
&ObjectOptions::default(),
)
.await
);
let stream = meta_obj.stream;
let mut data = vec![];
pin_mut!(stream);
while let Some(x) = stream.next().await {
let x = try_!(x);
data.put_slice(&x[..]);
}
let mut meta = try_!(BucketMetadata::unmarshal_from(&data[..]));
meta.tagging = None;
let data = try_!(meta.marshal_msg());
let len = data.len();
try_!(
store
.put_object(
RUSTFS_META_BUCKET,
BucketMetadata::new(bucket.as_str()).save_file_path().as_str(),
PutObjReader::new(StreamingBlob::from(Body::from(data)), len),
&ObjectOptions::default(),
)
.await
);
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn put_object_tagging(&self, req: S3Request<PutObjectTaggingInput>) -> S3Result<S3Response<PutObjectTaggingOutput>> {
let PutObjectTaggingInput {
bucket,
key: object,
tagging,
..
} = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let mut object_info = try_!(store.get_object_info(&bucket, &object, &ObjectOptions::default()).await);
object_info.tags = Some(tagging.tag_set.into_iter().map(|Tag { key, value }| (key, value)).collect());
try_!(
store
.put_object_info(&bucket, &object, object_info, &ObjectOptions::default())
.await
);
Ok(S3Response::new(PutObjectTaggingOutput { version_id: None }))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn get_object_tagging(&self, req: S3Request<GetObjectTaggingInput>) -> S3Result<S3Response<GetObjectTaggingOutput>> {
let GetObjectTaggingInput { bucket, key: object, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let object_info = try_!(store.get_object_info(&bucket, &object, &ObjectOptions::default()).await);
Ok(S3Response::new(GetObjectTaggingOutput {
tag_set: object_info
.tags
.map(|tags| tags.into_iter().map(|(key, value)| Tag { key, value }).collect())
.unwrap_or_else(|| vec![]),
version_id: None,
}))
}
#[tracing::instrument(level = "debug", skip(self))]
async fn delete_object_tagging(
&self,
req: S3Request<DeleteObjectTaggingInput>,
) -> S3Result<S3Response<DeleteObjectTaggingOutput>> {
let DeleteObjectTaggingInput { bucket, key: object, .. } = req.input;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = lock
.as_ref()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let mut object_info = try_!(store.get_object_info(&bucket, &object, &ObjectOptions::default()).await);
object_info.tags = None;
try_!(
store
.put_object_info(&bucket, &object, object_info, &ObjectOptions::default())
.await
);
Ok(S3Response::new(DeleteObjectTaggingOutput { version_id: None }))
}
}
#[allow(dead_code)]