mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 00:38:16 +00:00
merge main
This commit is contained in:
+67
-16
@@ -4,10 +4,6 @@ use super::options::extract_metadata;
|
||||
use super::options::put_opts;
|
||||
use crate::auth::get_condition_values;
|
||||
use crate::storage::access::ReqInfo;
|
||||
use crate::storage::error::to_s3_error;
|
||||
use crate::storage::options::copy_dst_opts;
|
||||
use crate::storage::options::copy_src_opts;
|
||||
use crate::storage::options::{extract_metadata_from_mime, get_opts};
|
||||
use bytes::Bytes;
|
||||
use common::error::Result;
|
||||
use ecstore::bucket::error::BucketMetadataError;
|
||||
@@ -44,7 +40,12 @@ use futures::pin_mut;
|
||||
use futures::{Stream, StreamExt};
|
||||
use http::HeaderMap;
|
||||
use lazy_static::lazy_static;
|
||||
use policy::auth;
|
||||
use policy::policy::action::Action;
|
||||
use policy::policy::action::S3Action;
|
||||
use policy::policy::BucketPolicy;
|
||||
use policy::policy::BucketPolicyArgs;
|
||||
use policy::policy::Validator;
|
||||
use s3s::dto::*;
|
||||
use s3s::s3_error;
|
||||
use s3s::S3Error;
|
||||
@@ -56,10 +57,18 @@ use std::fmt::Debug;
|
||||
use std::str::FromStr;
|
||||
use tokio_util::io::ReaderStream;
|
||||
use tokio_util::io::StreamReader;
|
||||
use tracing::{debug, error, info, warn};
|
||||
use tracing::debug;
|
||||
use tracing::error;
|
||||
use tracing::info;
|
||||
use tracing::warn;
|
||||
use transform_stream::AsyncTryStream;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::storage::error::to_s3_error;
|
||||
use crate::storage::options::copy_dst_opts;
|
||||
use crate::storage::options::copy_src_opts;
|
||||
use crate::storage::options::{extract_metadata_from_mime, get_opts};
|
||||
|
||||
macro_rules! try_ {
|
||||
($result:expr) => {
|
||||
match $result {
|
||||
@@ -554,7 +563,7 @@ impl S3 for FS {
|
||||
|
||||
let mut req = req;
|
||||
|
||||
if authorize_request(&mut req, vec![Action::S3Action(S3Action::ListAllMyBucketsAction)])
|
||||
if authorize_request(&mut req, Action::S3Action(S3Action::ListAllMyBucketsAction))
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
@@ -564,10 +573,10 @@ impl S3 for FS {
|
||||
req_info.bucket = Some(info.name.clone());
|
||||
|
||||
futures::executor::block_on(async {
|
||||
authorize_request(&mut req, vec![Action::S3Action(S3Action::ListBucketAction)])
|
||||
authorize_request(&mut req, Action::S3Action(S3Action::ListBucketAction))
|
||||
.await
|
||||
.is_ok()
|
||||
|| authorize_request(&mut req, vec![Action::S3Action(S3Action::GetBucketLocationAction)])
|
||||
|| authorize_request(&mut req, Action::S3Action(S3Action::GetBucketLocationAction))
|
||||
.await
|
||||
.is_ok()
|
||||
})
|
||||
@@ -1214,9 +1223,52 @@ impl S3 for FS {
|
||||
|
||||
async fn get_bucket_policy_status(
|
||||
&self,
|
||||
_req: S3Request<GetBucketPolicyStatusInput>,
|
||||
req: S3Request<GetBucketPolicyStatusInput>,
|
||||
) -> S3Result<S3Response<GetBucketPolicyStatusOutput>> {
|
||||
Err(s3_error!(NotImplemented, "GetBucketPolicyStatus is not implemented yet"))
|
||||
let GetBucketPolicyStatusInput { bucket, .. } = req.input;
|
||||
|
||||
let Some(store) = new_object_layer_fn() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
};
|
||||
|
||||
store
|
||||
.get_bucket_info(&bucket, &BucketOptions::default())
|
||||
.await
|
||||
.map_err(to_s3_error)?;
|
||||
|
||||
let conditions = get_condition_values(&req.headers, &auth::Credentials::default());
|
||||
|
||||
let read_olny = PolicySys::is_allowed(&BucketPolicyArgs {
|
||||
bucket: &bucket,
|
||||
action: Action::S3Action(S3Action::ListBucketAction),
|
||||
is_owner: false,
|
||||
account: "",
|
||||
groups: &None,
|
||||
conditions: &conditions,
|
||||
object: "",
|
||||
})
|
||||
.await;
|
||||
|
||||
let write_olny = PolicySys::is_allowed(&BucketPolicyArgs {
|
||||
bucket: &bucket,
|
||||
action: Action::S3Action(S3Action::PutObjectAction),
|
||||
is_owner: false,
|
||||
account: "",
|
||||
groups: &None,
|
||||
conditions: &conditions,
|
||||
object: "",
|
||||
})
|
||||
.await;
|
||||
|
||||
let is_public = read_olny && write_olny;
|
||||
|
||||
let output = GetBucketPolicyStatusOutput {
|
||||
policy_status: Some(PolicyStatus {
|
||||
is_public: Some(is_public),
|
||||
}),
|
||||
};
|
||||
|
||||
Ok(S3Response::new(output))
|
||||
}
|
||||
|
||||
async fn get_bucket_policy(&self, req: S3Request<GetBucketPolicyInput>) -> S3Result<S3Response<GetBucketPolicyOutput>> {
|
||||
@@ -1260,18 +1312,17 @@ impl S3 for FS {
|
||||
|
||||
// warn!("input policy {}", &policy);
|
||||
|
||||
let cfg = BucketPolicy::unmarshal(policy.as_bytes()).map_err(to_s3_error)?;
|
||||
let cfg: BucketPolicy =
|
||||
serde_json::from_str(&policy).map_err(|e| s3_error!(InvalidArgument, "parse policy faild {:?}", e))?;
|
||||
|
||||
// warn!("parse policy {:?}", &cfg);
|
||||
|
||||
if let Err(err) = cfg.validate(&bucket) {
|
||||
if let Err(err) = cfg.is_valid() {
|
||||
warn!("put_bucket_policy err input {:?}, {:?}", &policy, err);
|
||||
return Err(s3_error!(InvalidPolicyDocument));
|
||||
}
|
||||
|
||||
let data = cfg.marshal_msg().map_err(to_s3_error)?;
|
||||
let data = serde_json::to_vec(&cfg).map_err(|e| s3_error!(InternalError, "parse policy faild {:?}", e))?;
|
||||
|
||||
metadata_sys::update(&bucket, BUCKET_POLICY_CONFIG, data.into())
|
||||
metadata_sys::update(&bucket, BUCKET_POLICY_CONFIG, data)
|
||||
.await
|
||||
.map_err(to_s3_error)?;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user