separate signer.

fix ilm feature.
This commit is contained in:
likewu
2025-06-28 11:02:37 +08:00
parent e4690c48b4
commit 4ed84a6bc4
54 changed files with 498 additions and 1200 deletions
@@ -25,6 +25,8 @@ use tracing::{error, info};
use uuid::Uuid;
use xxhash_rust::xxh64;
//use rustfs_notify::{BucketNotificationConfig, Event, EventName, LogLevel, NotificationError, init_logger};
//use rustfs_notify::{initialize, notification_system};
use super::bucket_lifecycle_audit::{LcAuditEvent, LcEventSrc};
use super::lifecycle::{self, ExpirationOptions, IlmAction, Lifecycle, TransitionOptions};
use super::tier_last_day_stats::{DailyAllTierStats, LastDayTierStats};
@@ -109,6 +111,7 @@ struct ExpiryStats {
workers: AtomicI64,
}
#[allow(dead_code)]
impl ExpiryStats {
pub fn missed_tasks(&self) -> i64 {
self.missed_expiry_tasks.load(Ordering::SeqCst)
@@ -567,32 +570,6 @@ impl TransitionState {
}
}
struct AuditTierOp {
tier: String,
time_to_responsens: i64,
output_bytes: i64,
error: String,
}
impl AuditTierOp {
#[allow(clippy::new_ret_no_self)]
pub async fn new() -> Result<Self, std::io::Error> {
Ok(Self {
tier: String::from("tier"),
time_to_responsens: 0,
output_bytes: 0,
error: String::from(""),
})
}
pub fn string(&self) -> String {
format!(
"tier:{},respNS:{},tx:{},err:{}",
self.tier, self.time_to_responsens, self.output_bytes, self.error
)
}
}
pub async fn init_background_expiry(api: Arc<ECStore>) {
let mut workers = num_cpus::get() / 2;
//globalILMConfig.getExpirationWorkers()
@@ -717,6 +694,16 @@ pub async fn expire_transitioned_object(
host: GLOBAL_LocalNodeName.to_string(),
..Default::default()
});
/*let system = match notification_system() {
Some(sys) => sys,
None => {
let config = Config::new();
initialize(config).await?;
notification_system().expect("Failed to initialize notification system")
}
};
let event = Arc::new(Event::new_test_event("my-bucket", "document.pdf", EventName::ObjectCreatedPut));
system.send_event(event).await;*/
Ok(dobj)
}
@@ -845,4 +832,4 @@ pub struct RestoreObjectRequest {
pub output_location: OutputLocation,
}
const MAX_RESTORE_OBJECT_REQUEST_SIZE: i64 = 2 << 20;
const _MAX_RESTORE_OBJECT_REQUEST_SIZE: i64 = 2 << 20;
+1 -9
View File
@@ -25,7 +25,7 @@ pub const TRANSITION_PENDING: &str = "pending";
const ERR_LIFECYCLE_TOO_MANY_RULES: &str = "Lifecycle configuration allows a maximum of 1000 rules";
const ERR_LIFECYCLE_NO_RULE: &str = "Lifecycle configuration should have at least one rule";
const ERR_LIFECYCLE_DUPLICATE_ID: &str = "Rule ID must be unique. Found same ID for more than one rule";
const ERR_XML_NOT_WELL_FORMED: &str = "The XML you provided was not well-formed or did not validate against our published schema";
const _ERR_XML_NOT_WELL_FORMED: &str = "The XML you provided was not well-formed or did not validate against our published schema";
const ERR_LIFECYCLE_BUCKET_LOCKED: &str =
"ExpiredObjectAllVersions element and DelMarkerExpiration action cannot be used on an object locked bucket";
@@ -698,14 +698,6 @@ pub struct ExpirationOptions {
}
impl ExpirationOptions {
fn marshal_msg(&self, b: &[u8]) -> Result<Vec<u8>, std::io::Error> {
todo!();
}
fn unmarshal_msg(&self, bts: &[u8]) -> Result<Vec<u8>, std::io::Error> {
todo!();
}
fn msg_size(&self) -> i64 {
1 + 7 + 10
}
+3 -3
View File
@@ -7,11 +7,11 @@
use s3s::dto::{LifecycleRuleFilter, Transition};
const ERR_TRANSITION_INVALID_DAYS: &str = "Days must be 0 or greater when used with Transition";
const ERR_TRANSITION_INVALID_DATE: &str = "Date must be provided in ISO 8601 format";
const _ERR_TRANSITION_INVALID_DAYS: &str = "Days must be 0 or greater when used with Transition";
const _ERR_TRANSITION_INVALID_DATE: &str = "Date must be provided in ISO 8601 format";
const ERR_TRANSITION_INVALID: &str =
"Exactly one of Days (0 or greater) or Date (positive ISO 8601 format) should be present in Transition.";
const ERR_TRANSITION_DATE_NOT_MIDNIGHT: &str = "'Date' must be at midnight GMT";
const _ERR_TRANSITION_DATE_NOT_MIDNIGHT: &str = "'Date' must be at midnight GMT";
pub trait Filter {
fn test_tags(&self, user_tags: &str) -> bool;
@@ -65,6 +65,7 @@ impl LastDayTierStats {
}
}
#[allow(dead_code)]
fn merge(&self, m: LastDayTierStats) -> LastDayTierStats {
let mut cl = self.clone();
let mut cm = m.clone();
@@ -17,6 +17,7 @@ use crate::global::GLOBAL_TierConfigMgr;
static XXHASH_SEED: u64 = 0;
#[derive(Default)]
#[allow(dead_code)]
struct ObjSweeper {
object: String,
bucket: String,
@@ -29,6 +30,7 @@ struct ObjSweeper {
remote_object: String,
}
#[allow(dead_code)]
impl ObjSweeper {
#[allow(clippy::new_ret_no_self)]
pub async fn new(bucket: &str, object: &str) -> Result<Self, std::io::Error> {
@@ -103,6 +105,7 @@ impl ObjSweeper {
}
#[derive(Debug, Clone)]
#[allow(unused_assignments)]
pub struct Jentry {
obj_name: String,
version_id: String,
+9 -15
View File
@@ -4,21 +4,15 @@ use time::{OffsetDateTime, format_description};
use s3s::dto::{Date, ObjectLockLegalHold, ObjectLockLegalHoldStatus, ObjectLockRetention, ObjectLockRetentionMode};
use s3s::header::{X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE};
//const AMZ_OBJECTLOCK_BYPASS_RET_GOVERNANCE: &str = "X-Amz-Bypass-Governance-Retention";
//const AMZ_OBJECTLOCK_RETAIN_UNTIL_DATE: &str = "X-Amz-Object-Lock-Retain-Until-Date";
//const AMZ_OBJECTLOCK_MODE: &str = "X-Amz-Object-Lock-Mode";
//const AMZ_OBJECTLOCK_LEGALHOLD: &str = "X-Amz-Object-Lock-Legal-Hold";
// Commented out unused constants to avoid dead code warnings
// const ERR_MALFORMED_BUCKET_OBJECT_CONFIG: &str = "invalid bucket object lock config";
// const ERR_INVALID_RETENTION_DATE: &str = "date must be provided in ISO 8601 format";
// const ERR_PAST_OBJECTLOCK_RETAIN_DATE: &str = "the retain until date must be in the future";
// const ERR_UNKNOWN_WORMMODE_DIRECTIVE: &str = "unknown WORM mode directive";
// const ERR_OBJECTLOCK_MISSING_CONTENT_MD5: &str =
// "content-MD5 HTTP header is required for Put Object requests with Object Lock parameters";
// const ERR_OBJECTLOCK_INVALID_HEADERS: &str =
// "x-amz-object-lock-retain-until-date and x-amz-object-lock-mode must both be supplied";
// const ERR_MALFORMED_XML: &str = "the XML you provided was not well-formed or did not validate against our published schema";
const _ERR_MALFORMED_BUCKET_OBJECT_CONFIG: &str = "invalid bucket object lock config";
const _ERR_INVALID_RETENTION_DATE: &str = "date must be provided in ISO 8601 format";
const _ERR_PAST_OBJECTLOCK_RETAIN_DATE: &str = "the retain until date must be in the future";
const _ERR_UNKNOWN_WORMMODE_DIRECTIVE: &str = "unknown WORM mode directive";
const _ERR_OBJECTLOCK_MISSING_CONTENT_MD5: &str =
"content-MD5 HTTP header is required for Put Object requests with Object Lock parameters";
const _ERR_OBJECTLOCK_INVALID_HEADERS: &str =
"x-amz-object-lock-retain-until-date and x-amz-object-lock-mode must both be supplied";
const _ERR_MALFORMED_XML: &str = "the XML you provided was not well-formed or did not validate against our published schema";
pub fn utc_now_ntp() -> OffsetDateTime {
OffsetDateTime::now_utc()
@@ -17,7 +17,7 @@ impl BucketObjectLockSys {
}
pub async fn get(bucket: &str) -> Option<DefaultRetention> {
if let Some(object_lock_config) = get_object_lock_config(bucket).await {
if let Ok(object_lock_config) = get_object_lock_config(bucket).await {
if let Some(object_lock_rule) = object_lock_config.0.rule {
return object_lock_rule.default_retention;
}
+2 -1
View File
@@ -11,7 +11,7 @@ use std::collections::HashMap;
use crate::client::{api_put_object::PutObjectOptions, api_s3_datatypes::ObjectPart};
use crate::{disk::DiskAPI, store_api::GetObjectReader};
use reader::hasher::{Hasher, Sha256};
use rustfs_utils::hasher::{Hasher, Sha256};
use rustfs_utils::crypto::{base64_decode, base64_encode};
use s3s::header::{
X_AMZ_CHECKSUM_ALGORITHM, X_AMZ_CHECKSUM_CRC32, X_AMZ_CHECKSUM_CRC32C, X_AMZ_CHECKSUM_SHA1, X_AMZ_CHECKSUM_SHA256,
@@ -234,6 +234,7 @@ pub struct Checksum {
computed: bool,
}
#[allow(dead_code)]
impl Checksum {
fn new(t: ChecksumMode, b: &[u8]) -> Checksum {
if t.is_set() && b.len() == t.raw_byte_len() {
+9 -9
View File
@@ -1,10 +1,10 @@
use http::status::StatusCode;
use std::fmt::{self, Display, Formatter};
#[derive(Default, thiserror::Error, Debug, PartialEq)]
#[derive(Default, thiserror::Error, Debug, Clone, PartialEq)]
pub struct AdminError {
pub code: &'static str,
pub message: &'static str,
pub code: String,
pub message: String,
pub status_code: StatusCode,
}
@@ -15,18 +15,18 @@ impl Display for AdminError {
}
impl AdminError {
pub fn new(code: &'static str, message: &'static str, status_code: StatusCode) -> Self {
pub fn new(code: &str, message: &str, status_code: StatusCode) -> Self {
Self {
code,
message,
code: code.to_string(),
message: message.to_string(),
status_code,
}
}
pub fn msg(message: &'static str) -> Self {
pub fn msg(message: &str) -> Self {
Self {
code: "InternalError",
message,
code: "InternalError".to_string(),
message: message.to_string(),
status_code: StatusCode::INTERNAL_SERVER_ERROR,
}
}
+2
View File
@@ -61,6 +61,7 @@ impl TransitionClient {
}
#[derive(Default)]
#[allow(dead_code)]
pub struct GetRequest {
pub buffer: Vec<u8>,
pub offset: i64,
@@ -72,6 +73,7 @@ pub struct GetRequest {
pub setting_object_info: bool,
}
#[allow(dead_code)]
pub struct GetResponse {
pub size: i64,
//pub error: error,
+1 -2
View File
@@ -14,6 +14,7 @@ use tracing::warn;
use crate::client::api_error_response::err_invalid_argument;
#[derive(Default)]
#[allow(dead_code)]
pub struct AdvancedGetOptions {
replication_deletemarker: bool,
is_replication_ready_for_deletemarker: bool,
@@ -30,8 +31,6 @@ pub struct GetObjectOptions {
pub internal: AdvancedGetOptions,
}
type StatObjectOptions = GetObjectOptions;
impl Default for GetObjectOptions {
fn default() -> Self {
Self {
+1
View File
@@ -258,6 +258,7 @@ impl TransitionClient {
}
}
#[allow(dead_code)]
pub struct ListObjectsOptions {
reverse_versions: bool,
with_versions: bool,
+2 -1
View File
@@ -12,7 +12,7 @@ use std::{collections::HashMap, sync::Arc};
use time::{Duration, OffsetDateTime, macros::format_description};
use tracing::{error, info, warn};
use reader::hasher::Hasher;
use rustfs_utils::hasher::Hasher;
use s3s::dto::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode, ReplicationStatus};
use s3s::header::{
X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, X_AMZ_REPLICATION_STATUS,
@@ -124,6 +124,7 @@ impl Default for PutObjectOptions {
}
}
#[allow(dead_code)]
impl PutObjectOptions {
fn set_matche_tag(&mut self, etag: &str) {
if etag == "*" {
@@ -14,8 +14,6 @@ use crate::client::{
transition_api::TransitionClient,
};
const NULL_VERSION_ID: &str = "null";
pub fn is_object(reader: &ReaderImpl) -> bool {
todo!();
}
@@ -18,7 +18,7 @@ use tracing::{error, info};
use url::form_urlencoded::Serializer;
use uuid::Uuid;
use reader::hasher::Hasher;
use rustfs_utils::hasher::Hasher;
use s3s::header::{X_AMZ_EXPIRATION, X_AMZ_VERSION_ID};
use s3s::{Body, dto::StreamingBlob};
//use crate::disk::{Reader, BufferReader};
@@ -27,7 +27,7 @@ use crate::client::{
constants::ISO8601_DATEFORMAT,
transition_api::{ReaderImpl, RequestMetadata, TransitionClient, UploadInfo},
};
use reader::hasher::Hasher;
use rustfs_utils::hasher::Hasher;
use rustfs_utils::{crypto::base64_encode, path::trim_etag};
use s3s::header::{X_AMZ_EXPIRATION, X_AMZ_VERSION_ID};
+5 -2
View File
@@ -24,14 +24,15 @@ use crate::{
disk::DiskAPI,
store_api::{GetObjectReader, ObjectInfo, StorageAPI},
};
use reader::hasher::{sum_md5_base64, sum_sha256_hex};
use rustfs_utils::hasher::{sum_md5_base64, sum_sha256_hex};
use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH;
pub struct RemoveBucketOptions {
forced_elete: bool,
_forced_elete: bool,
}
#[derive(Debug)]
#[allow(dead_code)]
pub struct AdvancedRemoveOptions {
replication_delete_marker: bool,
replication_status: ReplicationStatus,
@@ -426,8 +427,10 @@ impl TransitionClient {
}
#[derive(Debug, Default)]
#[allow(dead_code)]
pub struct RemoveObjectError {
object_name: String,
#[allow(dead_code)]
version_id: String,
err: Option<std::io::Error>,
}
+1 -6
View File
@@ -43,6 +43,7 @@ pub struct ListBucketV2Result {
pub start_after: String,
}
#[allow(dead_code)]
pub struct Version {
etag: String,
is_latest: bool,
@@ -72,12 +73,6 @@ pub struct ListVersionsResult {
next_version_id_marker: String,
}
impl ListVersionsResult {
fn unmarshal_xml() -> Result<(), std::io::Error> {
todo!();
}
}
pub struct ListBucketResult {
common_prefixes: Vec<CommonPrefix>,
contents: Vec<transition_api::ObjectInfo>,
+3 -4
View File
@@ -17,8 +17,7 @@ use crate::client::{
api_error_response::{http_resp_to_error_response, to_error_response},
transition_api::{Document, TransitionClient},
};
use crate::signer;
use reader::hasher::{Hasher, Sha256};
use rustfs_utils::hasher::{Hasher, Sha256};
use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH;
use s3s::Body;
use s3s::S3ErrorCode;
@@ -151,7 +150,7 @@ impl TransitionClient {
}
if signer_type == SignatureType::SignatureV2 {
let req_builder = signer::sign_v2(req_builder, 0, &access_key_id, &secret_access_key, is_virtual_style);
let req_builder = rustfs_signer::sign_v2(req_builder, 0, &access_key_id, &secret_access_key, is_virtual_style);
let req = match req_builder.body(Body::empty()) {
Ok(req) => return Ok(req),
Err(err) => {
@@ -169,7 +168,7 @@ impl TransitionClient {
.headers_mut()
.expect("err")
.insert("X-Amz-Content-Sha256", content_sha256.parse().unwrap());
let req_builder = signer::sign_v4(req_builder, 0, &access_key_id, &secret_access_key, &session_token, "us-east-1");
let req_builder = rustfs_signer::sign_v4(req_builder, 0, &access_key_id, &secret_access_key, &session_token, "us-east-1");
let req = match req_builder.body(Body::empty()) {
Ok(req) => return Ok(req),
Err(err) => {
+1 -2
View File
@@ -1,4 +1,3 @@
#![allow(clippy::map_entry)]
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
@@ -28,4 +27,4 @@ pub const ISO8601_DATEFORMAT: &[FormatItem<'_>] =
pub const GET_OBJECT_ATTRIBUTES_TAGS: &str = "ETag,Checksum,StorageClass,ObjectSize,ObjectParts";
pub const GET_OBJECT_ATTRIBUTES_MAX_PARTS: i64 = 1000;
const RUSTFS_BUCKET_SOURCE_MTIME: &str = "X-RustFs-Source-Mtime";
const _RUSTFS_BUCKET_SOURCE_MTIME: &str = "X-RustFs-Source-Mtime";
-22
View File
@@ -141,28 +141,6 @@ impl ErrorResponse {
}
}
struct Error {
code: String,
message: String,
bucket_name: String,
key: String,
resource: String,
request_id: String,
host_id: String,
region: String,
server: String,
status_code: i64,
}
impl Display for Error {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
if self.message == "" {
return write!(f, "{}", format!("Error response code {}.", self.code));
}
write!(f, "{}", self.message)
}
}
pub fn xml_decoder<T>(body: &[u8]) -> Result<T, std::io::Error> {
todo!();
}
+1 -5
View File
@@ -24,16 +24,12 @@ pub struct PutObjReader {
//pub sealMD5Fn: SealMD5CurrFn,
}
#[allow(dead_code)]
impl PutObjReader {
pub fn new(raw_reader: HashReader) -> Self {
todo!();
}
fn size(&self) -> usize {
//self.reader.size()
todo!();
}
fn md5_current_hex_string(&self) -> String {
todo!();
}
+7 -25
View File
@@ -45,9 +45,8 @@ use crate::client::{
constants::{UNSIGNED_PAYLOAD, UNSIGNED_PAYLOAD_TRAILER},
credentials::{CredContext, Credentials, SignatureType, Static},
};
use crate::signer;
use crate::{checksum::ChecksumMode, store_api::GetObjectReader};
use reader::hasher::{MD5, Sha256};
use rustfs_utils::hasher::{MD5, Sha256};
use rustfs_rio::HashReader;
use rustfs_utils::{
net::get_endpoint_url,
@@ -57,7 +56,7 @@ use s3s::S3ErrorCode;
use s3s::dto::ReplicationStatus;
use s3s::{Body, dto::Owner};
const C_USER_AGENT_PREFIX: &str = "RustFS (linux; x86)";
const _C_USER_AGENT_PREFIX: &str = "RustFS (linux; x86)";
const C_USER_AGENT: &str = "RustFS (linux; x86)";
const SUCCESS_STATUS: [StatusCode; 3] = [StatusCode::OK, StatusCode::NO_CONTENT, StatusCode::PARTIAL_CONTENT];
@@ -433,9 +432,9 @@ impl TransitionClient {
}
if signer_type == SignatureType::SignatureV2 {
req_builder =
signer::pre_sign_v2(req_builder, &access_key_id, &secret_access_key, metadata.expires, is_virtual_host);
rustfs_signer::pre_sign_v2(req_builder, &access_key_id, &secret_access_key, metadata.expires, is_virtual_host);
} else if signer_type == SignatureType::SignatureV4 {
req_builder = signer::pre_sign_v4(
req_builder = rustfs_signer::pre_sign_v4(
req_builder,
&access_key_id,
&secret_access_key,
@@ -486,7 +485,7 @@ impl TransitionClient {
if signer_type == SignatureType::SignatureV2 {
req_builder =
signer::sign_v2(req_builder, metadata.content_length, &access_key_id, &secret_access_key, is_virtual_host);
rustfs_signer::sign_v2(req_builder, metadata.content_length, &access_key_id, &secret_access_key, is_virtual_host);
} else if metadata.stream_sha256 && !self.secure {
if metadata.trailer.len() > 0 {
//req.Trailer = metadata.trailer;
@@ -494,7 +493,7 @@ impl TransitionClient {
req_builder = req_builder.header(http::header::TRAILER, v.clone());
}
}
//req_builder = signer::streaming_sign_v4(req_builder, &access_key_id,
//req_builder = rustfs_signer::streaming_sign_v4(req_builder, &access_key_id,
// &secret_access_key, &session_token, &location, metadata.content_length, OffsetDateTime::now_utc(), self.sha256_hasher());
} else {
let mut sha_header = UNSIGNED_PAYLOAD.to_string();
@@ -509,7 +508,7 @@ impl TransitionClient {
req_builder = req_builder
.header::<HeaderName, HeaderValue>("X-Amz-Content-Sha256".parse().unwrap(), sha_header.parse().expect("err"));
req_builder = signer::sign_v4_trailer(
req_builder = rustfs_signer::sign_v4_trailer(
req_builder,
&access_key_id,
&secret_access_key,
@@ -614,23 +613,6 @@ impl TransitionClient {
}
}
struct LockedRandSource {
src: u64, //rand.Source,
}
impl LockedRandSource {
fn int63(&self) -> i64 {
/*let n = self.src.int63();
n*/
todo!();
}
fn seed(&self, seed: i64) {
//self.src.seed(seed);
todo!();
}
}
pub struct RequestMetadata {
pub pre_sign_url: bool,
pub bucket_name: String,
+4 -2
View File
@@ -810,6 +810,7 @@ pub fn error_resp_to_object_err(err: ErrorResponse, params: Vec<&str>) -> std::i
}
let r_err = err;
let err;
let bucket = bucket.to_string();
let object = object.to_string();
let version_id = version_id.to_string();
@@ -870,9 +871,10 @@ pub fn error_resp_to_object_err(err: ErrorResponse, params: Vec<&str>) -> std::i
/*S3ErrorCode::ReplicationPermissionCheck => {
err = std::io::Error::other(StorageError::ReplicationPermissionCheck);
}*/
_ => std::io::Error::other("err"),
_ => {
err = err_;
}
}
}
pub fn storage_to_object_err(err: Error, params: Vec<&str>) -> S3Error {
let storage_err = &err;
+1
View File
@@ -146,6 +146,7 @@ impl TierStats {
}
}
#[allow(dead_code)]
#[allow(dead_code)]
struct AllTierStats {
tiers: HashMap<String, TierStats>,
-1
View File
@@ -34,7 +34,6 @@ pub mod checksum;
pub mod client;
pub mod event;
pub mod event_notification;
pub mod signer;
pub mod tier;
pub use global::new_object_layer_fn;
-13
View File
@@ -1,13 +0,0 @@
pub mod ordered_qs;
pub mod request_signature_streaming;
pub mod request_signature_streaming_unsigned_trailer;
pub mod request_signature_v2;
pub mod request_signature_v4;
pub mod utils;
pub use request_signature_streaming::streaming_sign_v4;
pub use request_signature_v2::pre_sign_v2;
pub use request_signature_v2::sign_v2;
pub use request_signature_v4::pre_sign_v4;
pub use request_signature_v4::sign_v4;
pub use request_signature_v4::sign_v4_trailer;
-109
View File
@@ -1,109 +0,0 @@
//! Ordered query strings
use crate::signer::utils::stable_sort_by_first;
/// Immutable query string container
#[derive(Debug, Default, Clone)]
pub struct OrderedQs {
/// Ascending query strings
qs: Vec<(String, String)>,
}
/// [`OrderedQs`]
#[derive(Debug, thiserror::Error)]
#[error("ParseOrderedQsError: {inner}")]
pub struct ParseOrderedQsError {
/// url decode error
inner: serde_urlencoded::de::Error,
}
impl OrderedQs {
/// Constructs [`OrderedQs`] from vec
///
/// + strings must be url-decoded
#[cfg(test)]
#[must_use]
pub fn from_vec_unchecked(mut v: Vec<(String, String)>) -> Self {
stable_sort_by_first(&mut v);
Self { qs: v }
}
/// Parses [`OrderedQs`] from query
///
/// # Errors
/// Returns [`ParseOrderedQsError`] if query cannot be decoded
pub fn parse(query: &str) -> Result<Self, ParseOrderedQsError> {
let result = serde_urlencoded::from_str::<Vec<(String, String)>>(query);
let mut v = result.map_err(|e| ParseOrderedQsError { inner: e })?;
stable_sort_by_first(&mut v);
Ok(Self { qs: v })
}
#[must_use]
pub fn has(&self, name: &str) -> bool {
self.qs.binary_search_by_key(&name, |x| x.0.as_str()).is_ok()
}
/// Gets query values by name. Time `O(logn)`
pub fn get_all(&self, name: &str) -> impl Iterator<Item = &str> + use<'_> {
let qs = self.qs.as_slice();
let lower_bound = qs.partition_point(|x| x.0.as_str() < name);
let upper_bound = qs.partition_point(|x| x.0.as_str() <= name);
qs[lower_bound..upper_bound].iter().map(|x| x.1.as_str())
}
pub fn get_unique(&self, name: &str) -> Option<&str> {
let qs = self.qs.as_slice();
let lower_bound = qs.partition_point(|x| x.0.as_str() < name);
let mut iter = qs[lower_bound..].iter();
let pair = iter.next()?;
if let Some(following) = iter.next() {
if following.0 == name {
return None;
}
}
(pair.0.as_str() == name).then_some(pair.1.as_str())
}
}
impl AsRef<[(String, String)]> for OrderedQs {
fn as_ref(&self) -> &[(String, String)] {
self.qs.as_ref()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn tag() {
{
let query = "tagging";
let qs = OrderedQs::parse(query).unwrap();
assert_eq!(qs.as_ref(), &[("tagging".to_owned(), String::new())]);
assert_eq!(qs.get_unique("taggin"), None);
assert_eq!(qs.get_unique("tagging"), Some(""));
assert_eq!(qs.get_unique("taggingg"), None);
}
{
let query = "tagging&tagging";
let qs = OrderedQs::parse(query).unwrap();
assert_eq!(
qs.as_ref(),
&[("tagging".to_owned(), String::new()), ("tagging".to_owned(), String::new())]
);
assert_eq!(qs.get_unique("taggin"), None);
assert_eq!(qs.get_unique("tagging"), None);
assert_eq!(qs.get_unique("taggingg"), None);
}
}
}
@@ -1,74 +0,0 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use http::request::{self, Request};
use lazy_static::lazy_static;
use std::collections::HashMap;
use time::{OffsetDateTime, macros::format_description};
use super::request_signature_v4::{SERVICE_TYPE_S3, get_scope, get_signature, get_signing_key};
use rustfs_utils::{
crypto::{hex, hex_sha256, hex_sha256_chunk, hmac_sha256},
hash::EMPTY_STRING_SHA256_HASH,
};
const STREAMING_SIGN_ALGORITHM: &str = "STREAMING-AWS4-HMAC-SHA256-PAYLOAD";
const STREAMING_SIGN_TRAILER_ALGORITHM: &str = "STREAMING-AWS4-HMAC-SHA256-PAYLOAD-TRAILER";
const STREAMING_PAYLOAD_HDR: &str = "AWS4-HMAC-SHA256-PAYLOAD";
const STREAMING_TRAILER_HDR: &str = "AWS4-HMAC-SHA256-TRAILER";
const PAYLOAD_CHUNK_SIZE: i64 = 64 * 1024;
const CHUNK_SIGCONST_LEN: i64 = 17;
const SIGNATURESTR_LEN: i64 = 64;
const CRLF_LEN: i64 = 2;
const TRAILER_KV_SEPARATOR: &str = ":";
const TRAILER_SIGNATURE: &str = "x-amz-trailer-signature";
lazy_static! {
static ref ignored_streaming_headers: HashMap<String, bool> = {
let mut m = <HashMap<String, bool>>::new();
m.insert("authorization".to_string(), true);
m.insert("user-agent".to_string(), true);
m.insert("content-type".to_string(), true);
m
};
}
fn build_chunk_string_to_sign(t: OffsetDateTime, region: &str, previous_sig: &str, chunk_check_sum: &str) -> String {
let mut string_to_sign_parts = <Vec<String>>::new();
string_to_sign_parts.push(STREAMING_PAYLOAD_HDR.to_string());
let format = format_description!("[year][month][day]T[hour][minute][second]Z");
string_to_sign_parts.push(t.format(&format).unwrap());
string_to_sign_parts.push(get_scope(region, t, SERVICE_TYPE_S3));
string_to_sign_parts.push(previous_sig.to_string());
string_to_sign_parts.push(EMPTY_STRING_SHA256_HASH.to_string());
string_to_sign_parts.push(chunk_check_sum.to_string());
string_to_sign_parts.join("\n")
}
fn build_chunk_signature(
chunk_check_sum: &str,
req_time: OffsetDateTime,
region: &str,
previous_signature: &str,
secret_access_key: &str,
) -> String {
let chunk_string_to_sign = build_chunk_string_to_sign(req_time, region, previous_signature, chunk_check_sum);
let signing_key = get_signing_key(secret_access_key, region, req_time, SERVICE_TYPE_S3);
get_signature(signing_key, &chunk_string_to_sign)
}
pub fn streaming_sign_v4(
req: request::Builder,
access_key_id: &str,
secret_access_key: &str,
session_token: &str,
region: &str,
data_len: i64,
req_time: OffsetDateTime, /*, sh256: md5simd.Hasher*/
) -> request::Builder {
todo!();
}
@@ -1,17 +0,0 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
use http::request;
use time::OffsetDateTime;
pub fn streaming_unsigned_v4(
mut req: request::Builder,
session_token: &str,
data_len: i64,
req_time: OffsetDateTime,
) -> request::Builder {
todo!();
}
-238
View File
@@ -1,238 +0,0 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use bytes::{Bytes, BytesMut};
use http::request;
use hyper::Uri;
use std::collections::HashMap;
use std::fmt::Write;
use time::{OffsetDateTime, format_description, macros::format_description};
use rustfs_utils::crypto::{base64_encode, hex, hmac_sha1};
use super::utils::get_host_addr;
const SIGN_V4_ALGORITHM: &str = "AWS4-HMAC-SHA256";
const SIGN_V2_ALGORITHM: &str = "AWS";
fn encode_url2path(req: &request::Builder, virtual_host: bool) -> String {
let mut path = "".to_string();
//path = serde_urlencoded::to_string(req.uri_ref().unwrap().path().unwrap()).unwrap();
path = req.uri_ref().unwrap().path().to_string();
path
}
pub fn pre_sign_v2(
mut req: request::Builder,
access_key_id: &str,
secret_access_key: &str,
expires: i64,
virtual_host: bool,
) -> request::Builder {
if access_key_id == "" || secret_access_key == "" {
return req;
}
let d = OffsetDateTime::now_utc();
let d = d.replace_time(time::Time::from_hms(0, 0, 0).unwrap());
let epoch_expires = d.unix_timestamp() + expires;
let mut headers = req.headers_mut().expect("err");
let expires_str = headers.get("Expires");
if expires_str.is_none() {
headers.insert("Expires", format!("{:010}", epoch_expires).parse().unwrap());
}
let string_to_sign = pre_string_to_sign_v2(&req, virtual_host);
let signature = hex(hmac_sha1(secret_access_key, string_to_sign));
let result = serde_urlencoded::from_str::<HashMap<String, String>>(req.uri_ref().unwrap().query().unwrap());
let mut query = result.unwrap_or_default();
if get_host_addr(&req).contains(".storage.googleapis.com") {
query.insert("GoogleAccessId".to_string(), access_key_id.to_string());
} else {
query.insert("AWSAccessKeyId".to_string(), access_key_id.to_string());
}
query.insert("Expires".to_string(), format!("{:010}", epoch_expires));
let uri = req.uri_ref().unwrap().clone();
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(
format!("{}?{}&Signature={}", uri.path(), serde_urlencoded::to_string(&query).unwrap(), signature)
.parse()
.unwrap(),
);
let req = req.uri(Uri::from_parts(parts).unwrap());
req
}
fn post_pre_sign_signature_v2(policy_base64: &str, secret_access_key: &str) -> String {
let signature = hex(hmac_sha1(secret_access_key, policy_base64));
signature
}
pub fn sign_v2(
mut req: request::Builder,
content_len: i64,
access_key_id: &str,
secret_access_key: &str,
virtual_host: bool,
) -> request::Builder {
if access_key_id == "" || secret_access_key == "" {
return req;
}
let d = OffsetDateTime::now_utc();
let d2 = d.replace_time(time::Time::from_hms(0, 0, 0).unwrap());
let string_to_sign = string_to_sign_v2(&req, virtual_host);
let mut headers = req.headers_mut().expect("err");
let date = headers.get("Date").unwrap();
if date.to_str().unwrap() == "" {
headers.insert(
"Date",
d2.format(&format_description::well_known::Rfc2822)
.unwrap()
.to_string()
.parse()
.unwrap(),
);
}
let mut auth_header = format!("{} {}:", SIGN_V2_ALGORITHM, access_key_id);
let auth_header = format!("{}{}", auth_header, base64_encode(&hmac_sha1(secret_access_key, string_to_sign)));
headers.insert("Authorization", auth_header.parse().unwrap());
req
}
fn pre_string_to_sign_v2(req: &request::Builder, virtual_host: bool) -> String {
let mut buf = BytesMut::new();
write_pre_sign_v2_headers(&mut buf, &req);
write_canonicalized_headers(&mut buf, &req);
write_canonicalized_resource(&mut buf, &req, virtual_host);
String::from_utf8(buf.to_vec()).unwrap()
}
fn write_pre_sign_v2_headers(buf: &mut BytesMut, req: &request::Builder) {
let _ = buf.write_str(req.method_ref().unwrap().as_str());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Content-Md5").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Content-Type").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Expires").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
}
fn string_to_sign_v2(req: &request::Builder, virtual_host: bool) -> String {
let mut buf = BytesMut::new();
write_sign_v2_headers(&mut buf, &req);
write_canonicalized_headers(&mut buf, &req);
write_canonicalized_resource(&mut buf, &req, virtual_host);
String::from_utf8(buf.to_vec()).unwrap()
}
fn write_sign_v2_headers(buf: &mut BytesMut, req: &request::Builder) {
let _ = buf.write_str(req.method_ref().unwrap().as_str());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Content-Md5").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Content-Type").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
let _ = buf.write_str(req.headers_ref().unwrap().get("Date").unwrap().to_str().unwrap());
let _ = buf.write_char('\n');
}
fn write_canonicalized_headers(buf: &mut BytesMut, req: &request::Builder) {
let mut proto_headers = Vec::<String>::new();
let mut vals = HashMap::<String, Vec<String>>::new();
for k in req.headers_ref().expect("err").keys() {
let lk = k.as_str().to_lowercase();
if lk.starts_with("x-amz") {
proto_headers.push(lk.clone());
let vv = req
.headers_ref()
.expect("err")
.get_all(k)
.iter()
.map(|e| e.to_str().unwrap().to_string())
.collect();
vals.insert(lk, vv);
}
}
proto_headers.sort();
for k in proto_headers {
let _ = buf.write_str(&k);
let _ = buf.write_char(':');
for (idx, v) in vals[&k].iter().enumerate() {
if idx > 0 {
let _ = buf.write_char(',');
}
let _ = buf.write_str(v);
}
let _ = buf.write_char('\n');
}
}
const INCLUDED_QUERY: &[&str] = &[
"acl",
"delete",
"lifecycle",
"location",
"logging",
"notification",
"partNumber",
"policy",
"requestPayment",
"response-cache-control",
"response-content-disposition",
"response-content-encoding",
"response-content-language",
"response-content-type",
"response-expires",
"uploadId",
"uploads",
"versionId",
"versioning",
"versions",
"website",
];
fn write_canonicalized_resource(buf: &mut BytesMut, req: &request::Builder, virtual_host: bool) {
let request_url = req.uri_ref().unwrap();
let _ = buf.write_str(&encode_url2path(req, virtual_host));
if request_url.query().unwrap() != "" {
let mut n: i64 = 0;
let result = serde_urlencoded::from_str::<HashMap<String, Vec<String>>>(req.uri_ref().unwrap().query().unwrap());
let mut vals = result.unwrap_or_default();
for resource in INCLUDED_QUERY {
let vv = &vals[*resource];
if vv.len() > 0 {
n += 1;
match n {
1 => {
let _ = buf.write_char('?');
}
_ => {
let _ = buf.write_char('&');
let _ = buf.write_str(resource);
if vv[0].len() > 0 {
let _ = buf.write_char('=');
let _ = buf.write_str(&vv[0]);
}
}
}
}
}
}
}
-740
View File
@@ -1,740 +0,0 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use bytes::{Bytes, BytesMut};
use http::HeaderMap;
use http::Uri;
use http::header::TRAILER;
use http::request::{self, Request};
use lazy_static::lazy_static;
use std::collections::HashMap;
use std::fmt::Write;
use time::{OffsetDateTime, format_description, macros::datetime, macros::format_description};
use tracing::{debug, error, info, warn};
use super::ordered_qs::OrderedQs;
use super::request_signature_streaming_unsigned_trailer::streaming_unsigned_v4;
use super::utils::stable_sort_by_first;
use super::utils::{get_host_addr, sign_v4_trim_all};
use crate::client::constants::UNSIGNED_PAYLOAD;
use rustfs_utils::crypto::{hex, hex_sha256, hmac_sha256};
use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH;
pub const SIGN_V4_ALGORITHM: &str = "AWS4-HMAC-SHA256";
pub const SERVICE_TYPE_S3: &str = "s3";
pub const SERVICE_TYPE_STS: &str = "sts";
lazy_static! {
static ref v4_ignored_headers: HashMap<String, bool> = {
let mut m = <HashMap<String, bool>>::new();
m.insert("accept-encoding".to_string(), true);
m.insert("authorization".to_string(), true);
m.insert("user-agent".to_string(), true);
m
};
}
pub fn get_signing_key(secret: &str, loc: &str, t: OffsetDateTime, service_type: &str) -> [u8; 32] {
let mut s = "AWS4".to_string();
s.push_str(secret);
let format = format_description!("[year][month][day]");
let date = hmac_sha256(s.into_bytes(), t.format(&format).unwrap().into_bytes());
let location = hmac_sha256(date, loc);
let service = hmac_sha256(location, service_type);
let signing_key = hmac_sha256(service, "aws4_request");
signing_key
}
pub fn get_signature(signing_key: [u8; 32], string_to_sign: &str) -> String {
hex(hmac_sha256(signing_key, string_to_sign))
}
pub fn get_scope(location: &str, t: OffsetDateTime, service_type: &str) -> String {
let format = format_description!("[year][month][day]");
let mut ans = String::from("");
ans.push_str(&t.format(&format).unwrap().to_string());
ans.push('/');
ans.push_str(location); // TODO: use a `Region` type
ans.push('/');
ans.push_str(service_type);
ans.push_str("/aws4_request");
ans
}
fn get_credential(access_key_id: &str, location: &str, t: OffsetDateTime, service_type: &str) -> String {
let scope = get_scope(location, t, service_type);
let mut s = access_key_id.to_string();
s.push_str("/");
s.push_str(&scope);
s
}
fn get_hashed_payload(req: &request::Builder) -> String {
let headers = req.headers_ref().unwrap();
let mut hashed_payload = "";
if let Some(payload) = headers.get("X-Amz-Content-Sha256") {
hashed_payload = payload.to_str().unwrap();
}
if hashed_payload == "" {
hashed_payload = UNSIGNED_PAYLOAD;
}
hashed_payload.to_string()
}
fn get_canonical_headers(req: &request::Builder, ignored_headers: &HashMap<String, bool>) -> String {
let mut headers = Vec::<String>::new();
let mut vals = HashMap::<String, Vec<String>>::new();
for k in req.headers_ref().expect("err").keys() {
if ignored_headers.get(&k.to_string()).is_some() {
continue;
}
headers.push(k.as_str().to_lowercase());
let vv = req
.headers_ref()
.expect("err")
.get_all(k)
.iter()
.map(|e| e.to_str().unwrap().to_string())
.collect();
vals.insert(k.as_str().to_lowercase(), vv);
}
if !header_exists("host", &headers) {
headers.push("host".to_string());
}
headers.sort();
debug!("get_canonical_headers vals: {:?}", vals);
debug!("get_canonical_headers headers: {:?}", headers);
let mut buf = BytesMut::new();
for k in headers {
let _ = buf.write_str(&k);
let _ = buf.write_char(':');
let k: &str = &k;
match k {
"host" => {
let _ = buf.write_str(&get_host_addr(&req));
let _ = buf.write_char('\n');
}
_ => {
for (idx, v) in vals[k].iter().enumerate() {
if idx > 0 {
let _ = buf.write_char(',');
}
let _ = buf.write_str(&sign_v4_trim_all(v));
}
let _ = buf.write_char('\n');
}
}
}
String::from_utf8(buf.to_vec()).unwrap()
}
fn header_exists(key: &str, headers: &[String]) -> bool {
for k in headers {
if k == key {
return true;
}
}
false
}
fn get_signed_headers(req: &request::Builder, ignored_headers: &HashMap<String, bool>) -> String {
let mut headers = Vec::<String>::new();
let headers_ref = req.headers_ref().expect("err");
debug!("get_signed_headers headers: {:?}", headers_ref);
for (k, _) in headers_ref {
if ignored_headers.get(&k.to_string()).is_some() {
continue;
}
headers.push(k.as_str().to_lowercase());
}
if !header_exists("host", &headers) {
headers.push("host".to_string());
}
headers.sort();
headers.join(";")
}
fn get_canonical_request(req: &request::Builder, ignored_headers: &HashMap<String, bool>, hashed_payload: &str) -> String {
let mut canonical_query_string = "".to_string();
if let Some(q) = req.uri_ref().unwrap().query() {
// Parse query string into key-value pairs
let mut query_params: Vec<(String, String)> = Vec::new();
for param in q.split('&') {
if let Some((key, value)) = param.split_once('=') {
query_params.push((key.to_string(), value.to_string()));
} else {
query_params.push((param.to_string(), "".to_string()));
}
}
// Sort by key name
query_params.sort_by(|a, b| a.0.cmp(&b.0));
// Build canonical query string
let sorted_params: Vec<String> = query_params
.iter()
.map(|(k, v)| if v.is_empty() { k.clone() } else { format!("{}={}", k, v) })
.collect();
canonical_query_string = sorted_params.join("&");
canonical_query_string = canonical_query_string.replace("+", "%20");
}
let mut canonical_request = <Vec<String>>::new();
canonical_request.push(req.method_ref().unwrap().to_string());
canonical_request.push(req.uri_ref().unwrap().path().to_string());
canonical_request.push(canonical_query_string);
canonical_request.push(get_canonical_headers(&req, ignored_headers));
canonical_request.push(get_signed_headers(&req, ignored_headers));
canonical_request.push(hashed_payload.to_string());
canonical_request.join("\n")
}
fn get_string_to_sign_v4(t: OffsetDateTime, location: &str, canonical_request: &str, service_type: &str) -> String {
let mut string_to_sign = SIGN_V4_ALGORITHM.to_string();
string_to_sign.push('\n');
let format = format_description!("[year][month][day]T[hour][minute][second]Z");
string_to_sign.push_str(&t.format(&format).unwrap());
string_to_sign.push('\n');
string_to_sign.push_str(&get_scope(location, t, service_type));
string_to_sign.push('\n');
string_to_sign.push_str(&hex_sha256(canonical_request.as_bytes(), |s| s.to_string()));
string_to_sign
}
pub fn pre_sign_v4(
req: request::Builder,
access_key_id: &str,
secret_access_key: &str,
session_token: &str,
location: &str,
expires: i64,
t: OffsetDateTime,
) -> request::Builder {
if access_key_id == "" || secret_access_key == "" {
return req;
}
//let t = OffsetDateTime::now_utc();
//let date = AmzDate::parse(timestamp).unwrap();
let t2 = t.replace_time(time::Time::from_hms(0, 0, 0).unwrap());
//let credential = get_scope(location, t, SERVICE_TYPE_S3);
let credential = get_credential(access_key_id, location, t, SERVICE_TYPE_S3);
let signed_headers = get_signed_headers(&req, &v4_ignored_headers);
let mut query = <Vec<(String, String)>>::new();
if let Some(q) = req.uri_ref().unwrap().query() {
let result = serde_urlencoded::from_str::<Vec<(String, String)>>(q);
query = result.unwrap_or_default();
}
query.push(("X-Amz-Algorithm".to_string(), SIGN_V4_ALGORITHM.to_string()));
let format = format_description!("[year][month][day]T[hour][minute][second]Z");
query.push(("X-Amz-Date".to_string(), t.format(&format).unwrap().to_string()));
query.push(("X-Amz-Expires".to_string(), format!("{:010}", expires)));
query.push(("X-Amz-SignedHeaders".to_string(), signed_headers));
query.push(("X-Amz-Credential".to_string(), credential));
if session_token != "" {
query.push(("X-Amz-Security-Token".to_string(), session_token.to_string()));
}
let uri = req.uri_ref().unwrap().clone();
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(
format!("{}?{}", uri.path(), serde_urlencoded::to_string(&query).unwrap())
.parse()
.unwrap(),
);
let req = req.uri(Uri::from_parts(parts).unwrap());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
let string_to_sign = get_string_to_sign_v4(t, location, &canonical_request, SERVICE_TYPE_S3);
//println!("canonical_request: \n{}\n", canonical_request);
//println!("string_to_sign: \n{}\n", string_to_sign);
let signing_key = get_signing_key(secret_access_key, location, t, SERVICE_TYPE_S3);
let signature = get_signature(signing_key, &string_to_sign);
let uri = req.uri_ref().unwrap().clone();
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(
format!(
"{}?{}&X-Amz-Signature={}",
uri.path(),
serde_urlencoded::to_string(&query).unwrap(),
signature
)
.parse()
.unwrap(),
);
let req = req.uri(Uri::from_parts(parts).unwrap());
req
}
fn post_pre_sign_signature_v4(policy_base64: &str, t: OffsetDateTime, secret_access_key: &str, location: &str) -> String {
let signing_key = get_signing_key(secret_access_key, location, t, SERVICE_TYPE_S3);
let signature = get_signature(signing_key, policy_base64);
signature
}
fn sign_v4_sts(mut req: request::Builder, access_key_id: &str, secret_access_key: &str, location: &str) -> request::Builder {
sign_v4_inner(req, 0, access_key_id, secret_access_key, "", location, SERVICE_TYPE_STS, HeaderMap::new())
}
fn sign_v4_inner(
mut req: request::Builder,
content_len: i64,
access_key_id: &str,
secret_access_key: &str,
session_token: &str,
location: &str,
service_type: &str,
trailer: HeaderMap,
) -> request::Builder {
if access_key_id == "" || secret_access_key == "" {
return req;
}
let t = OffsetDateTime::now_utc();
let t2 = t.replace_time(time::Time::from_hms(0, 0, 0).unwrap());
let mut headers = req.headers_mut().expect("err");
let format = format_description!("[year][month][day]T[hour][minute][second]Z");
headers.insert("X-Amz-Date", t.format(&format).unwrap().to_string().parse().unwrap());
if session_token != "" {
headers.insert("X-Amz-Security-Token", session_token.parse().unwrap());
}
if trailer.len() > 0 {
for (k, _) in &trailer {
headers.append("X-Amz-Trailer", k.as_str().to_lowercase().parse().unwrap());
}
headers.insert("Content-Encoding", "aws-chunked".parse().unwrap());
headers.insert("x-amz-decoded-content-length", format!("{:010}", content_len).parse().unwrap());
}
if service_type == SERVICE_TYPE_STS {
headers.remove("X-Amz-Content-Sha256");
}
let hashed_payload = get_hashed_payload(&req);
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &hashed_payload);
let string_to_sign = get_string_to_sign_v4(t, location, &canonical_request, service_type);
let signing_key = get_signing_key(secret_access_key, location, t, service_type);
let credential = get_credential(access_key_id, location, t2, service_type);
let signed_headers = get_signed_headers(&req, &v4_ignored_headers);
let signature = get_signature(signing_key, &string_to_sign);
//debug!("\n\ncanonical_request: \n{}\nstring_to_sign: \n{}\nsignature: \n{}\n\n", &canonical_request, &string_to_sign, &signature);
let mut headers = req.headers_mut().expect("err");
let auth = format!(
"{} Credential={}, SignedHeaders={}, Signature={}",
SIGN_V4_ALGORITHM, credential, signed_headers, signature
);
headers.insert("Authorization", auth.parse().unwrap());
if trailer.len() > 0 {
//req.Trailer = trailer;
for (_, v) in &trailer {
headers.append(http::header::TRAILER, v.clone());
}
return streaming_unsigned_v4(req, &session_token, content_len, t);
}
req
}
fn unsigned_trailer(mut req: request::Builder, content_len: i64, trailer: HeaderMap) {
if trailer.len() > 0 {
return;
}
let t = OffsetDateTime::now_utc();
let t = t.replace_time(time::Time::from_hms(0, 0, 0).unwrap());
let mut headers = req.headers_mut().expect("err");
let format = format_description!("[year][month][day]T[hour][minute][second]Z");
headers.insert("X-Amz-Date", t.format(&format).unwrap().to_string().parse().unwrap());
for (k, _) in &trailer {
headers.append("X-Amz-Trailer", k.as_str().to_lowercase().parse().unwrap());
}
headers.insert("Content-Encoding", "aws-chunked".parse().unwrap());
headers.insert("x-amz-decoded-content-length", format!("{:010}", content_len).parse().unwrap());
if trailer.len() > 0 {
for (_, v) in &trailer {
headers.append(http::header::TRAILER, v.clone());
}
}
streaming_unsigned_v4(req, "", content_len, t);
}
pub fn sign_v4(
mut req: request::Builder,
content_len: i64,
access_key_id: &str,
secret_access_key: &str,
session_token: &str,
location: &str,
) -> request::Builder {
sign_v4_inner(
req,
content_len,
access_key_id,
secret_access_key,
session_token,
location,
SERVICE_TYPE_S3,
HeaderMap::new(),
)
}
pub fn sign_v4_trailer(
req: request::Builder,
access_key_id: &str,
secret_access_key: &str,
session_token: &str,
location: &str,
trailer: HeaderMap,
) -> request::Builder {
sign_v4_inner(
req,
0,
access_key_id,
secret_access_key,
session_token,
location,
SERVICE_TYPE_S3,
trailer,
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn example_list_objects() {
// let access_key_id = "AKIAIOSFODNN7EXAMPLE";
let secret_access_key = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
let timestamp = "20130524T000000Z";
let t = datetime!(2013-05-24 0:00 UTC);
// let bucket = "examplebucket";
let region = "us-east-1";
let service = "s3";
let path = "/";
let mut req = Request::builder()
.method(http::Method::GET)
.uri("http://examplebucket.s3.amazonaws.com/?");
let mut headers = req.headers_mut().expect("err");
headers.insert("host", "examplebucket.s3.amazonaws.com".parse().unwrap());
headers.insert(
"x-amz-content-sha256",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.parse()
.unwrap(),
);
headers.insert("x-amz-date", timestamp.parse().unwrap());
let mut query = <Vec<(String, String)>>::new();
query.push(("max-keys".to_string(), "2".to_string()));
query.push(("prefix".to_string(), "J".to_string()));
let uri = req.uri_ref().unwrap().clone();
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(
format!("{}?{}", uri.path(), serde_urlencoded::to_string(&query).unwrap())
.parse()
.unwrap(),
);
let req = req.uri(Uri::from_parts(parts).unwrap());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
assert_eq!(
canonical_request,
concat!(
"GET\n",
"/\n",
"max-keys=2&prefix=J\n",
"host:examplebucket.s3.amazonaws.com\n",
"x-amz-content-sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855\n",
"x-amz-date:",
"20130524T000000Z",
"\n",
"\n",
"host;x-amz-content-sha256;x-amz-date\n",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
)
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
assert_eq!(
string_to_sign,
concat!(
"AWS4-HMAC-SHA256\n",
"20130524T000000Z",
"\n",
"20130524/us-east-1/s3/aws4_request\n",
"df57d21db20da04d7fa30298dd4488ba3a2b47ca3a489c74750e0f1e7df1b9b7",
)
);
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
assert_eq!(signature, "34b48302e7b5fa45bde8084f4b7868a86f0a534bc59db6670ed5711ef69dc6f7");
}
#[test]
fn example_signature() {
// let access_key_id = "rustfsadmin";
let secret_access_key = "rustfsadmin";
let timestamp = "20250505T011054Z";
let t = datetime!(2025-05-05 01:10:54 UTC);
// let bucket = "mblock2";
let region = "us-east-1";
let service = "s3";
let path = "/mblock2/";
let mut req = Request::builder()
.method(http::Method::GET)
.uri("http://192.168.1.11:9020/mblock2/?");
let mut headers = req.headers_mut().expect("err");
headers.insert("host", "192.168.1.11:9020".parse().unwrap());
headers.insert(
"x-amz-content-sha256",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.parse()
.unwrap(),
);
headers.insert("x-amz-date", timestamp.parse().unwrap());
let mut query: Vec<(String, String)> = Vec::new();
let uri = req.uri_ref().unwrap().clone();
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(
format!("{}?{}", uri.path(), serde_urlencoded::to_string(&query).unwrap())
.parse()
.unwrap(),
);
let req = req.uri(Uri::from_parts(parts).unwrap());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
println!("canonical_request: \n{}\n", canonical_request);
assert_eq!(
canonical_request,
concat!(
"GET\n",
"/mblock2/\n",
"\n",
"host:192.168.1.11:9020\n",
"x-amz-content-sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855\n",
"x-amz-date:",
"20250505T011054Z",
"\n",
"\n",
"host;x-amz-content-sha256;x-amz-date\n",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
)
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
println!("string_to_sign: \n{}\n", string_to_sign);
assert_eq!(
string_to_sign,
concat!(
"AWS4-HMAC-SHA256\n",
"20250505T011054Z",
"\n",
"20250505/us-east-1/s3/aws4_request\n",
"c2960d00cc7de7bed3e2e2d1330ec298ded8f78a231c1d32dedac72ebec7f9b0",
)
);
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
println!("signature: \n{}\n", signature);
assert_eq!(signature, "df4116595e27b0dfd1103358947d9199378cc6386c4657abd8c5f0b11ebb4931");
}
#[test]
fn example_signature2() {
// let access_key_id = "rustfsadmin";
let secret_access_key = "rustfsadmin";
let timestamp = "20250507T051030Z";
let t = datetime!(2025-05-07 05:10:30 UTC);
// let bucket = "mblock2";
let region = "us-east-1";
let service = "s3";
let path = "/mblock2/";
let mut req = Request::builder().method(http::Method::GET).uri("http://192.168.1.11:9020/mblock2/?list-type=2&encoding-type=url&prefix=mypre&delimiter=%2F&fetch-owner=true&max-keys=1");
let mut headers = req.headers_mut().expect("err");
headers.insert("host", "192.168.1.11:9020".parse().unwrap());
headers.insert(
"x-amz-content-sha256",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.parse()
.unwrap(),
);
headers.insert("x-amz-date", timestamp.parse().unwrap());
/*let uri = req.uri_ref().unwrap().clone();
println!("{:?}", uri);
let mut canonical_query_string = "".to_string();
if let Some(q) = uri.query() {
let result = serde_urlencoded::from_str::<Vec<String>>(q);
let mut query = result.unwrap_or_default();
query.sort();
canonical_query_string = query.join("&");
canonical_query_string.replace("+", "%20");
}
let mut parts = req.uri_ref().unwrap().clone().into_parts();
parts.path_and_query = Some(format!("{}?{}", uri.path(), canonical_query_string).parse().unwrap());
let req = req.uri(Uri::from_parts(parts).unwrap());*/
println!("{:?}", req.uri_ref().unwrap().query());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
println!("canonical_request: \n{}\n", canonical_request);
assert_eq!(
canonical_request,
concat!(
"GET\n",
"/mblock2/\n",
"delimiter=%2F&encoding-type=url&fetch-owner=true&list-type=2&max-keys=1&prefix=mypre\n",
"host:192.168.1.11:9020\n",
"x-amz-content-sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855\n",
"x-amz-date:",
"20250507T051030Z",
"\n",
"\n",
"host;x-amz-content-sha256;x-amz-date\n",
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
)
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
println!("string_to_sign: \n{}\n", string_to_sign);
assert_eq!(
string_to_sign,
concat!(
"AWS4-HMAC-SHA256\n",
"20250507T051030Z",
"\n",
"20250507/us-east-1/s3/aws4_request\n",
"e6db9e09e9c873aff0b9ca170998b4753f6a6c36c90bc2dca80613affb47f999",
)
);
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
println!("signature: \n{}\n", signature);
assert_eq!(signature, "760278c9a77d5c245ac83d85917bddc3e3b14343091e8f4ad8edbbf73107d685");
}
#[test]
fn example_presigned_url() {
use hyper::Uri;
let access_key_id = "AKIAIOSFODNN7EXAMPLE";
let secret_access_key = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
let timestamp = "20130524T000000Z";
let t = datetime!(2013-05-24 0:00 UTC);
// let bucket = "mblock2";
let region = "us-east-1";
let service = "s3";
let path = "/";
let session_token = "";
let mut req = Request::builder()
.method(http::Method::GET)
.uri("http://examplebucket.s3.amazonaws.com/test.txt");
let mut headers = req.headers_mut().expect("err");
headers.insert("host", "examplebucket.s3.amazonaws.com".parse().unwrap());
req = pre_sign_v4(req, access_key_id, secret_access_key, "", region, 86400, t);
let mut canonical_request = req.method_ref().unwrap().as_str().to_string();
canonical_request.push('\n');
canonical_request.push_str(req.uri_ref().unwrap().path());
canonical_request.push('\n');
canonical_request.push_str(req.uri_ref().unwrap().query().unwrap());
canonical_request.push('\n');
canonical_request.push_str(&get_canonical_headers(&req, &v4_ignored_headers));
canonical_request.push('\n');
canonical_request.push_str(&get_signed_headers(&req, &v4_ignored_headers));
canonical_request.push('\n');
canonical_request.push_str(&get_hashed_payload(&req));
//println!("canonical_request: \n{}\n", canonical_request);
assert_eq!(
canonical_request,
concat!(
"GET\n",
"/test.txt\n",
"X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Date=20130524T000000Z&X-Amz-Expires=0000086400&X-Amz-SignedHeaders=host&X-Amz-Credential=AKIAIOSFODNN7EXAMPLE%2F20130524%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Signature=98f1c9f47b39a4c40662680a9b029b046b7da5542c2e35d67edb8ff18d2ccf5c\n",
"host:examplebucket.s3.amazonaws.com\n",
"\n",
"host\n",
"UNSIGNED-PAYLOAD",
)
);
}
#[test]
fn example_presigned_url2() {
use hyper::Uri;
let access_key_id = "rustfsadmin";
let secret_access_key = "rustfsadmin";
let timestamp = "20130524T000000Z";
let t = datetime!(2013-05-24 0:00 UTC);
// let bucket = "mblock2";
let region = "us-east-1";
let service = "s3";
let path = "/mblock2/";
let session_token = "";
let mut req = Request::builder().method(http::Method::GET).uri("http://192.168.1.11:9020/mblock2/test.txt?delimiter=%2F&fetch-owner=true&prefix=mypre&encoding-type=url&max-keys=1&list-type=2");
let mut headers = req.headers_mut().expect("err");
headers.insert("host", "192.168.1.11:9020".parse().unwrap());
req = pre_sign_v4(req, access_key_id, secret_access_key, "", region, 86400, t);
let mut canonical_request = req.method_ref().unwrap().as_str().to_string();
canonical_request.push('\n');
canonical_request.push_str(req.uri_ref().unwrap().path());
canonical_request.push('\n');
canonical_request.push_str(req.uri_ref().unwrap().query().unwrap());
canonical_request.push('\n');
canonical_request.push_str(&get_canonical_headers(&req, &v4_ignored_headers));
canonical_request.push('\n');
canonical_request.push_str(&get_signed_headers(&req, &v4_ignored_headers));
canonical_request.push('\n');
canonical_request.push_str(&get_hashed_payload(&req));
//println!("canonical_request: \n{}\n", canonical_request);
assert_eq!(
canonical_request,
concat!(
"GET\n",
"/mblock2/test.txt\n",
"delimiter=%2F&fetch-owner=true&prefix=mypre&encoding-type=url&max-keys=1&list-type=2&X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Date=20130524T000000Z&X-Amz-Expires=0000086400&X-Amz-SignedHeaders=host&X-Amz-Credential=rustfsadmin%2F20130524%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Signature=30e6b8c920512f0d12cba77a7c39612bff7f5f8148f4dc35cdd18f4b15a12477\n",
"host:192.168.1.11:9020\n",
"\n",
"host\n",
"UNSIGNED-PAYLOAD",
)
);
}
}
-33
View File
@@ -1,33 +0,0 @@
use http::request;
pub fn get_host_addr(req: &request::Builder) -> String {
let host = req.headers_ref().expect("err").get("host");
let uri = req.uri_ref().unwrap();
let req_host;
if let Some(port) = uri.port() {
req_host = format!("{}:{}", uri.host().unwrap(), port);
} else {
req_host = uri.host().unwrap().to_string();
}
if let Some(host) = host {
if req_host != *host.to_str().unwrap() {
return host.to_str().unwrap().to_string();
}
}
/*if req.uri_ref().unwrap().host().is_some() {
return req.uri_ref().unwrap().host().unwrap();
}*/
req_host
}
pub fn sign_v4_trim_all(input: &str) -> String {
let ss = input.split_whitespace().collect::<Vec<_>>();
ss.join(" ")
}
pub fn stable_sort_by_first<T>(v: &mut [(T, T)])
where
T: Ord,
{
v.sort_by(|lhs, rhs| lhs.0.cmp(&rhs.0));
}
+60 -60
View File
@@ -53,42 +53,40 @@ pub const TIER_CONFIG_FORMAT: u16 = 1;
pub const TIER_CONFIG_V1: u16 = 1;
pub const TIER_CONFIG_VERSION: u16 = 1;
const _TIER_CFG_REFRESH_AT_HDR: &str = "X-RustFS-TierCfg-RefreshedAt";
lazy_static! {
//pub static ref TIER_CONFIG_PATH: PathBuf = path_join(&[PathBuf::from(RUSTFS_CONFIG_PREFIX), PathBuf::from(TIER_CONFIG_FILE)]);
pub static ref ERR_TIER_MISSING_CREDENTIALS: AdminError = AdminError {
code: "XRustFSAdminTierMissingCredentials".to_string(),
message: "Specified remote credentials are empty".to_string(),
status_code: StatusCode::FORBIDDEN,
};
pub static ref ERR_TIER_BACKEND_IN_USE: AdminError = AdminError {
code: "XRustFSAdminTierBackendInUse".to_string(),
message: "Specified remote tier is already in use".to_string(),
status_code: StatusCode::CONFLICT,
};
pub static ref ERR_TIER_TYPE_UNSUPPORTED: AdminError = AdminError {
code: "XRustFSAdminTierTypeUnsupported".to_string(),
message: "Specified tier type is unsupported".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_BACKEND_NOT_EMPTY: AdminError = AdminError {
code: "XRustFSAdminTierBackendNotEmpty".to_string(),
message: "Specified remote backend is not empty".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_INVALID_CONFIG: AdminError = AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: "Unable to setup remote tier, check tier configuration".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
}
const TIER_CFG_REFRESH_AT_HDR: &str = "X-RustFS-TierCfg-RefreshedAt";
pub const ERR_TIER_MISSING_CREDENTIALS: AdminError = AdminError {
code: "XRustFSAdminTierMissingCredentials",
message: "Specified remote credentials are empty",
status_code: StatusCode::FORBIDDEN,
};
pub const ERR_TIER_BACKEND_IN_USE: AdminError = AdminError {
code: "XRustFSAdminTierBackendInUse",
message: "Specified remote tier is already in use",
status_code: StatusCode::CONFLICT,
};
pub const ERR_TIER_TYPE_UNSUPPORTED: AdminError = AdminError {
code: "XRustFSAdminTierTypeUnsupported",
message: "Specified tier type is unsupported",
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_BACKEND_NOT_EMPTY: AdminError = AdminError {
code: "XRustFSAdminTierBackendNotEmpty",
message: "Specified remote backend is not empty",
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_INVALID_CONFIG: AdminError = AdminError {
code: "XRustFSAdminTierInvalidConfig",
message: "Unable to setup remote tier, check tier configuration",
status_code: StatusCode::BAD_REQUEST,
};
#[derive(Serialize, Deserialize)]
pub struct TierConfigMgr {
#[serde(skip)]
@@ -108,20 +106,13 @@ impl TierConfigMgr {
pub fn unmarshal(data: &[u8]) -> std::result::Result<TierConfigMgr, std::io::Error> {
let cfg: TierConfigMgr = serde_json::from_slice(data)?;
//let mut cfg = TierConfigMgr(m);
//let mut cfg = m;
Ok(cfg)
}
pub fn marshal(&self) -> std::result::Result<Bytes, std::io::Error> {
let data = serde_json::to_vec(&self)?;
//let mut data = Vec<u8>::with_capacity(self.msg_size()+4);
let mut data = Bytes::from(data);
//LittleEndian::write_u16(&mut data[0..2], TIER_CONFIG_FORMAT);
//LittleEndian::write_u16(&mut data[2..4], TIER_CONFIG_VERSION);
Ok(data)
}
@@ -144,12 +135,12 @@ impl TierConfigMgr {
pub async fn add(&mut self, tier: TierConfig, force: bool) -> std::result::Result<(), AdminError> {
let tier_name = &tier.name;
if tier_name != tier_name.to_uppercase().as_str() {
return Err(ERR_TIER_NAME_NOT_UPPERCASE);
return Err(ERR_TIER_NAME_NOT_UPPERCASE.clone());
}
let (_, b) = self.is_tier_name_in_use(tier_name);
if b {
return Err(ERR_TIER_ALREADY_EXISTS);
return Err(ERR_TIER_ALREADY_EXISTS.clone());
}
let d = new_warm_backend(&tier, true).await?;
@@ -159,19 +150,22 @@ impl TierConfigMgr {
match in_use {
Ok(b) => {
if b {
return Err(ERR_TIER_BACKEND_IN_USE);
return Err(ERR_TIER_BACKEND_IN_USE.clone());
}
}
Err(err) => {
warn!("tier add failed, err: {:?}", err);
if err.to_string().contains("connect") {
return Err(ERR_TIER_CONNECT_ERR);
return Err(ERR_TIER_CONNECT_ERR.clone());
} else if err.to_string().contains("authorization") {
return Err(ERR_TIER_INVALID_CREDENTIALS);
return Err(ERR_TIER_INVALID_CREDENTIALS.clone());
} else if err.to_string().contains("bucket") {
return Err(ERR_TIER_BUCKET_NOT_FOUND);
return Err(ERR_TIER_BUCKET_NOT_FOUND.clone());
}
return Err(ERR_TIER_PERM_ERR);
let mut e = ERR_TIER_PERM_ERR.clone();
e.message.push('.');
e.message.push_str(&err.to_string());
return Err(e);
}
}
}
@@ -185,21 +179,21 @@ impl TierConfigMgr {
pub async fn remove(&mut self, tier_name: &str, force: bool) -> std::result::Result<(), AdminError> {
let d = self.get_driver(tier_name).await;
if let Err(err) = d {
match err {
ERR_TIER_NOT_FOUND => {
return Ok(());
}
_ => {
return Err(err);
}
if err.code == ERR_TIER_NOT_FOUND.code {
return Ok(());
} else {
return Err(err);
}
}
if !force {
let inuse = d.expect("err").in_use().await;
if let Err(err) = inuse {
return Err(ERR_TIER_PERM_ERR);
let mut e = ERR_TIER_PERM_ERR.clone();
e.message.push('.');
e.message.push_str(&err.to_string());
return Err(e);
} else if inuse.expect("err") {
return Err(ERR_TIER_BACKEND_NOT_EMPTY);
return Err(ERR_TIER_BACKEND_NOT_EMPTY.clone());
}
}
self.tiers.remove(tier_name);
@@ -254,7 +248,7 @@ impl TierConfigMgr {
pub async fn edit(&mut self, tier_name: &str, creds: TierCreds) -> std::result::Result<(), AdminError> {
let (tier_type, exists) = self.is_tier_name_in_use(tier_name);
if !exists {
return Err(ERR_TIER_NOT_FOUND);
return Err(ERR_TIER_NOT_FOUND.clone());
}
let mut cfg = self.tiers[tier_name].clone();
@@ -276,7 +270,7 @@ impl TierConfigMgr {
TierType::RustFS => {
let mut rustfs = cfg.rustfs.as_mut().expect("err");
if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS);
return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
rustfs.access_key = creds.access_key;
rustfs.secret_key = creds.secret_key;
@@ -284,7 +278,7 @@ impl TierConfigMgr {
TierType::MinIO => {
let mut minio = cfg.minio.as_mut().expect("err");
if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS);
return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
minio.access_key = creds.access_key;
minio.secret_key = creds.secret_key;
@@ -304,7 +298,7 @@ impl TierConfigMgr {
Entry::Vacant(e) => {
let t = self.tiers.get(tier_name);
if t.is_none() {
return Err(ERR_TIER_NOT_FOUND);
return Err(ERR_TIER_NOT_FOUND.clone());
}
let d = new_warm_backend(t.expect("err"), false).await?;
e.insert(d)
@@ -332,6 +326,12 @@ impl TierConfigMgr {
Ok(())
}
pub async fn clear_tier(&mut self, force: bool) -> std::result::Result<(), AdminError> {
self.tiers.clear();
self.driver_cache.clear();
Ok(())
}
#[tracing::instrument(level = "debug", name = "tier_save", skip(self))]
pub async fn save(&self) -> std::result::Result<(), std::io::Error> {
let Some(api) = new_object_layer_fn() else {
+11 -31
View File
@@ -1,10 +1,3 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use serde::{Deserialize, Serialize};
use std::fmt::Display;
use tracing::info;
@@ -12,9 +5,6 @@ use tracing::info;
const C_TIER_CONFIG_VER: &str = "v1";
const ERR_TIER_NAME_EMPTY: &str = "remote tier name empty";
const ERR_TIER_INVALID_CONFIG: &str = "invalid tier config";
const ERR_TIER_INVALID_CONFIG_VERSION: &str = "invalid tier config version";
const ERR_TIER_TYPE_UNSUPPORTED: &str = "unsupported tier type";
#[derive(Serialize, Deserialize, Default, Debug, Clone)]
pub enum TierType {
@@ -128,20 +118,8 @@ impl Clone for TierConfig {
}
}
#[allow(dead_code)]
impl TierConfig {
pub fn unmarshal(data: &[u8]) -> Result<TierConfig, std::io::Error> {
/*let m: HashMap<String, HashMap<String, KVS>> = serde_json::from_slice(data)?;
let mut cfg = TierConfig(m);
cfg.set_defaults();
Ok(cfg)*/
todo!();
}
pub fn marshal(&self) -> Result<Vec<u8>, std::io::Error> {
let data = serde_json::to_vec(&self)?;
Ok(data)
}
fn endpoint(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().expect("err").endpoint.clone(),
@@ -198,14 +176,14 @@ impl TierConfig {
pub struct TierS3 {
pub name: String,
pub endpoint: String,
#[serde(rename = "accesskey")]
#[serde(rename = "accessKey")]
pub access_key: String,
#[serde(rename = "secretkey")]
#[serde(rename = "secretKey")]
pub secret_key: String,
pub bucket: String,
pub prefix: String,
pub region: String,
#[serde(rename = "storageclass")]
#[serde(rename = "storageClass")]
pub storage_class: String,
#[serde(skip)]
pub aws_role: bool,
@@ -220,6 +198,7 @@ pub struct TierS3 {
}
impl TierS3 {
#[allow(dead_code)]
fn new<F>(name: &str, access_key: &str, secret_key: &str, bucket: &str, options: Vec<F>) -> Result<TierConfig, std::io::Error>
where
F: Fn(TierS3) -> Box<Result<(), std::io::Error>> + Send + Sync + 'static,
@@ -258,14 +237,14 @@ impl TierS3 {
pub struct TierRustFS {
pub name: String,
pub endpoint: String,
#[serde(rename = "accesskey")]
#[serde(rename = "accessKey")]
pub access_key: String,
#[serde(rename = "secretkey")]
#[serde(rename = "secretKey")]
pub secret_key: String,
pub bucket: String,
pub prefix: String,
pub region: String,
#[serde(rename = "storageclass")]
#[serde(rename = "storageClass")]
pub storage_class: String,
}
@@ -274,9 +253,9 @@ pub struct TierRustFS {
pub struct TierMinIO {
pub name: String,
pub endpoint: String,
#[serde(rename = "accesskey")]
#[serde(rename = "accessKey")]
pub access_key: String,
#[serde(rename = "secretkey")]
#[serde(rename = "secretKey")]
pub secret_key: String,
pub bucket: String,
pub prefix: String,
@@ -284,6 +263,7 @@ pub struct TierMinIO {
}
impl TierMinIO {
#[allow(dead_code)]
fn new<F>(
name: &str,
endpoint: &str,
+1 -23
View File
@@ -1,29 +1,7 @@
#![allow(unused_imports)]
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use crate::tier::tier::TierConfigMgr;
#[allow(dead_code)]
impl TierConfigMgr {
fn decode_msg(/*dc *msgp.Reader*/) -> Result<(), std::io::Error> {
todo!();
}
fn encode_msg(/*en *msgp.Writer*/) -> Result<(), std::io::Error> {
todo!();
}
pub fn marshal_msg(&self, b: &[u8]) -> Result<Vec<u8>, std::io::Error> {
todo!();
}
pub fn unmarshal_msg(buf: &[u8]) -> Result<Self, std::io::Error> {
todo!();
}
pub fn msg_size(&self) -> usize {
100
}
+43 -40
View File
@@ -1,50 +1,53 @@
use crate::client::admin_handler_utils::AdminError;
use http::status::StatusCode;
use lazy_static::lazy_static;
pub const ERR_TIER_ALREADY_EXISTS: AdminError = AdminError {
code: "XRustFSAdminTierAlreadyExists",
message: "Specified remote tier already exists",
status_code: StatusCode::CONFLICT,
};
lazy_static! {
pub static ref ERR_TIER_ALREADY_EXISTS: AdminError = AdminError {
code: "XRustFSAdminTierAlreadyExists".to_string(),
message: "Specified remote tier already exists".to_string(),
status_code: StatusCode::CONFLICT,
};
pub const ERR_TIER_NOT_FOUND: AdminError = AdminError {
code: "XRustFSAdminTierNotFound",
message: "Specified remote tier was not found",
status_code: StatusCode::NOT_FOUND,
};
pub static ref ERR_TIER_NOT_FOUND: AdminError = AdminError {
code: "XRustFSAdminTierNotFound".to_string(),
message: "Specified remote tier was not found".to_string(),
status_code: StatusCode::NOT_FOUND,
};
pub const ERR_TIER_NAME_NOT_UPPERCASE: AdminError = AdminError {
code: "XRustFSAdminTierNameNotUpperCase",
message: "Tier name must be in uppercase",
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_NAME_NOT_UPPERCASE: AdminError = AdminError {
code: "XRustFSAdminTierNameNotUpperCase".to_string(),
message: "Tier name must be in uppercase".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_BUCKET_NOT_FOUND: AdminError = AdminError {
code: "XRustFSAdminTierBucketNotFound",
message: "Remote tier bucket not found",
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_BUCKET_NOT_FOUND: AdminError = AdminError {
code: "XRustFSAdminTierBucketNotFound".to_string(),
message: "Remote tier bucket not found".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_INVALID_CREDENTIALS: AdminError = AdminError {
code: "XRustFSAdminTierInvalidCredentials",
message: "Invalid remote tier credentials",
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_INVALID_CREDENTIALS: AdminError = AdminError {
code: "XRustFSAdminTierInvalidCredentials".to_string(),
message: "Invalid remote tier credentials".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_RESERVED_NAME: AdminError = AdminError {
code: "XRustFSAdminTierReserved",
message: "Cannot use reserved tier name",
status_code: StatusCode::BAD_REQUEST,
};
pub static ref ERR_TIER_RESERVED_NAME: AdminError = AdminError {
code: "XRustFSAdminTierReserved".to_string(),
message: "Cannot use reserved tier name".to_string(),
status_code: StatusCode::BAD_REQUEST,
};
pub const ERR_TIER_PERM_ERR: AdminError = AdminError {
code: "TierPermErr",
message: "Tier Perm Err",
status_code: StatusCode::OK,
};
pub static ref ERR_TIER_PERM_ERR: AdminError = AdminError {
code: "TierPermErr".to_string(),
message: "Tier Perm Err".to_string(),
status_code: StatusCode::OK,
};
pub const ERR_TIER_CONNECT_ERR: AdminError = AdminError {
code: "TierConnectErr",
message: "Tier Connect Err",
status_code: StatusCode::OK,
};
pub static ref ERR_TIER_CONNECT_ERR: AdminError = AdminError {
code: "TierConnectErr".to_string(),
message: "Tier Connect Err".to_string(),
status_code: StatusCode::OK,
};
}
+22 -10
View File
@@ -7,14 +7,14 @@
use bytes::Bytes;
use std::collections::HashMap;
use http::StatusCode;
use crate::client::{
admin_handler_utils::AdminError,
transition_api::{ReadCloser, ReaderImpl},
};
use crate::error::is_err_bucket_not_found;
use crate::tier::{
tier::{ERR_TIER_INVALID_CONFIG, ERR_TIER_TYPE_UNSUPPORTED},
tier::ERR_TIER_TYPE_UNSUPPORTED,
tier_config::{TierConfig, TierType},
tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_PERM_ERR},
warm_backend_minio::WarmBackendMinIO,
@@ -54,7 +54,7 @@ pub async fn check_warm_backend(w: Option<&WarmBackendImpl>) -> Result<(), Admin
.put(PROBE_OBJECT, ReaderImpl::Body(Bytes::from("RustFS".as_bytes().to_vec())), 5)
.await;
if let Err(err) = remote_version_id {
return Err(ERR_TIER_PERM_ERR);
return Err(ERR_TIER_PERM_ERR.clone());
}
let r = w.get(PROBE_OBJECT, "", WarmBackendGetOpts::default()).await;
@@ -67,11 +67,11 @@ pub async fn check_warm_backend(w: Option<&WarmBackendImpl>) -> Result<(), Admin
return Err(ERR_TIER_MISSING_CREDENTIALS);
}*/
//else {
return Err(ERR_TIER_PERM_ERR);
return Err(ERR_TIER_PERM_ERR.clone());
//}
}
if let Err(err) = w.remove(PROBE_OBJECT, &remote_version_id.expect("err")).await {
return Err(ERR_TIER_PERM_ERR);
return Err(ERR_TIER_PERM_ERR.clone());
};
Ok(())
}
@@ -82,8 +82,12 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
TierType::S3 => {
let dd = WarmBackendS3::new(tier.s3.as_ref().expect("err"), &tier.name).await;
if let Err(err) = dd {
info!("{}", err);
return Err(ERR_TIER_INVALID_CONFIG);
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("err")));
}
@@ -91,7 +95,11 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
let dd = WarmBackendRustFS::new(tier.rustfs.as_ref().expect("err"), &tier.name).await;
if let Err(err) = dd {
warn!("{}", err);
return Err(ERR_TIER_INVALID_CONFIG);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("err")));
}
@@ -99,12 +107,16 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
let dd = WarmBackendMinIO::new(tier.minio.as_ref().expect("err"), &tier.name).await;
if let Err(err) = dd {
warn!("{}", err);
return Err(ERR_TIER_INVALID_CONFIG);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("err")));
}
_ => {
return Err(ERR_TIER_TYPE_UNSUPPORTED);
return Err(ERR_TIER_TYPE_UNSUPPORTED.clone());
}
}
+1 -1
View File
@@ -23,7 +23,7 @@ use tracing::warn;
const MAX_MULTIPART_PUT_OBJECT_SIZE: i64 = 1024 * 1024 * 1024 * 1024 * 5;
const MAX_PARTS_COUNT: i64 = 10000;
const MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
pub struct WarmBackendMinIO(WarmBackendS3);
+1 -1
View File
@@ -22,7 +22,7 @@ use crate::tier::{
const MAX_MULTIPART_PUT_OBJECT_SIZE: i64 = 1024 * 1024 * 1024 * 1024 * 5;
const MAX_PARTS_COUNT: i64 = 10000;
const MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
pub struct WarmBackendRustFS(WarmBackendS3);
-9
View File
@@ -95,15 +95,6 @@ impl WarmBackendS3 {
})
}
fn to_object_err(&self, err: ErrorResponse, params: Vec<&str>) -> std::io::Error {
let mut object = "";
if params.len() >= 1 {
object = params.first().cloned().unwrap_or_default();
}
error_resp_to_object_err(err, vec![&self.bucket, &self.get_dest(object)])
}
pub fn get_dest(&self, object: &str) -> String {
let mut dest_obj = object.to_string();
if self.prefix != "" {