mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-31 17:28:12 +00:00
fix(notify): validate bucket notification filter rules (#2806)
This commit is contained in:
@@ -91,6 +91,99 @@ fn to_internal_error(err: impl Display) -> S3Error {
|
|||||||
S3Error::with_message(S3ErrorCode::InternalError, format!("{err}"))
|
S3Error::with_message(S3ErrorCode::InternalError, format!("{err}"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn is_valid_notification_filter_value(value: &str) -> bool {
|
||||||
|
if value.len() > 1024 || value.contains('\\') {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
!value.split('/').any(|segment| segment == "." || segment == "..")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn invalid_filter_value_message(cfg_scope: &str, value: &str) -> String {
|
||||||
|
format!("invalid notification filter value (len={}) ({cfg_scope})", value.len())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn invalid_filter_name_message(cfg_scope: &str, name: &str) -> String {
|
||||||
|
format!(
|
||||||
|
"invalid notification filter name (len={}) (only 'prefix'/'suffix' are supported) ({cfg_scope})",
|
||||||
|
name.len()
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn validate_notification_filter_rules(
|
||||||
|
filter: Option<&NotificationConfigurationFilter>,
|
||||||
|
cfg_kind: &str,
|
||||||
|
cfg_id: Option<&str>,
|
||||||
|
) -> S3Result<()> {
|
||||||
|
let Some(filter) = filter else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let Some(s3key_filter) = filter.key.as_ref() else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let Some(rules) = s3key_filter.filter_rules.as_ref() else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
|
||||||
|
let mut has_prefix = false;
|
||||||
|
let mut has_suffix = false;
|
||||||
|
let cfg_scope = cfg_id.map_or_else(|| cfg_kind.to_string(), |id| format!("{cfg_kind} id={id}"));
|
||||||
|
|
||||||
|
for rule in rules {
|
||||||
|
let Some(name) = rule.name.as_ref() else {
|
||||||
|
return Err(s3_error!(InvalidArgument, "invalid notification filter rule: missing Name ({cfg_scope})"));
|
||||||
|
};
|
||||||
|
let Some(value) = rule.value.as_ref() else {
|
||||||
|
return Err(s3_error!(
|
||||||
|
InvalidArgument,
|
||||||
|
"invalid notification filter rule: missing Value ({cfg_scope})"
|
||||||
|
));
|
||||||
|
};
|
||||||
|
|
||||||
|
if !is_valid_notification_filter_value(value) {
|
||||||
|
return Err(s3_error!(InvalidArgument, "{}", invalid_filter_value_message(&cfg_scope, value)));
|
||||||
|
}
|
||||||
|
|
||||||
|
match name.as_str() {
|
||||||
|
"prefix" => {
|
||||||
|
if has_prefix {
|
||||||
|
return Err(s3_error!(InvalidArgument, "duplicate notification filter name 'prefix' ({cfg_scope})"));
|
||||||
|
}
|
||||||
|
has_prefix = true;
|
||||||
|
}
|
||||||
|
"suffix" => {
|
||||||
|
if has_suffix {
|
||||||
|
return Err(s3_error!(InvalidArgument, "duplicate notification filter name 'suffix' ({cfg_scope})"));
|
||||||
|
}
|
||||||
|
has_suffix = true;
|
||||||
|
}
|
||||||
|
other => {
|
||||||
|
return Err(s3_error!(InvalidArgument, "{}", invalid_filter_name_message(&cfg_scope, other)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn validate_notification_configuration_filters(notification_configuration: &NotificationConfiguration) -> S3Result<()> {
|
||||||
|
if let Some(queue_configs) = notification_configuration.queue_configurations.as_ref() {
|
||||||
|
for cfg in queue_configs {
|
||||||
|
validate_notification_filter_rules(cfg.filter.as_ref(), "QueueConfiguration", cfg.id.as_deref())?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if let Some(topic_configs) = notification_configuration.topic_configurations.as_ref() {
|
||||||
|
for cfg in topic_configs {
|
||||||
|
validate_notification_filter_rules(cfg.filter.as_ref(), "TopicConfiguration", cfg.id.as_deref())?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if let Some(lambda_configs) = notification_configuration.lambda_function_configurations.as_ref() {
|
||||||
|
for cfg in lambda_configs {
|
||||||
|
validate_notification_filter_rules(cfg.filter.as_ref(), "LambdaFunctionConfiguration", cfg.id.as_deref())?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
fn sr_bucket_meta_item(bucket: String, item_type: &str) -> SRBucketMeta {
|
fn sr_bucket_meta_item(bucket: String, item_type: &str) -> SRBucketMeta {
|
||||||
SRBucketMeta {
|
SRBucketMeta {
|
||||||
bucket,
|
bucket,
|
||||||
@@ -1497,6 +1590,8 @@ impl DefaultBucketUsecase {
|
|||||||
..
|
..
|
||||||
} = req.input;
|
} = req.input;
|
||||||
|
|
||||||
|
validate_notification_configuration_filters(¬ification_configuration)?;
|
||||||
|
|
||||||
let Some(store) = new_object_layer_fn() else {
|
let Some(store) = new_object_layer_fn() else {
|
||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||||
};
|
};
|
||||||
@@ -2879,6 +2974,163 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_invalid_filter_name() {
|
||||||
|
let raw_name = "Prefix".repeat(100);
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![FilterRule {
|
||||||
|
name: Some(FilterRuleName::from(raw_name.clone())),
|
||||||
|
value: Some("uploads/".to_string()),
|
||||||
|
}]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
let msg = err.message().unwrap_or_default();
|
||||||
|
assert!(msg.contains("len="), "error message should include summarized length");
|
||||||
|
assert!(!msg.contains(&raw_name), "error message should not echo full raw filter name");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_duplicate_prefix_rules() {
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![
|
||||||
|
FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)),
|
||||||
|
value: Some("uploads/".to_string()),
|
||||||
|
},
|
||||||
|
FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)),
|
||||||
|
value: Some("images/".to_string()),
|
||||||
|
},
|
||||||
|
]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_invalid_filter_value() {
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)),
|
||||||
|
value: Some("../secret".to_string()),
|
||||||
|
}]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
let msg = err.message().unwrap_or_default();
|
||||||
|
assert!(msg.contains("len="), "error message should include summarized length");
|
||||||
|
assert!(!msg.contains("../secret"), "error message should not echo full raw filter value");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_missing_filter_name() {
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![FilterRule {
|
||||||
|
name: None,
|
||||||
|
value: Some("uploads/".to_string()),
|
||||||
|
}]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_missing_filter_value() {
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::PREFIX)),
|
||||||
|
value: None,
|
||||||
|
}]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_notification_configuration_filters_rejects_duplicate_suffix_rules() {
|
||||||
|
let cfg = NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![
|
||||||
|
FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)),
|
||||||
|
value: Some(".csv".to_string()),
|
||||||
|
},
|
||||||
|
FilterRule {
|
||||||
|
name: Some(FilterRuleName::from_static(FilterRuleName::SUFFIX)),
|
||||||
|
value: Some(".log".to_string()),
|
||||||
|
},
|
||||||
|
]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = validate_notification_configuration_filters(&cfg).unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn execute_put_bucket_policy_returns_internal_error_when_store_uninitialized() {
|
async fn execute_put_bucket_policy_returns_internal_error_when_store_uninitialized() {
|
||||||
let input = PutBucketPolicyInput::builder()
|
let input = PutBucketPolicyInput::builder()
|
||||||
@@ -2894,6 +3146,36 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn execute_put_bucket_notification_configuration_rejects_invalid_filter_before_store_lookup() {
|
||||||
|
let input = PutBucketNotificationConfigurationInput::builder()
|
||||||
|
.bucket("test-bucket".to_string())
|
||||||
|
.notification_configuration(NotificationConfiguration {
|
||||||
|
queue_configurations: Some(vec![QueueConfiguration {
|
||||||
|
id: Some("q1".to_string()),
|
||||||
|
queue_arn: "arn:rustfs:sqs:us-east-1:1:webhook".to_string(),
|
||||||
|
events: vec!["s3:ObjectCreated:*".to_string().into()],
|
||||||
|
filter: Some(NotificationConfigurationFilter {
|
||||||
|
key: Some(S3KeyFilter {
|
||||||
|
filter_rules: Some(vec![FilterRule {
|
||||||
|
name: Some(FilterRuleName::from("Prefix".to_string())),
|
||||||
|
value: Some("uploads/".to_string()),
|
||||||
|
}]),
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
}]),
|
||||||
|
..Default::default()
|
||||||
|
})
|
||||||
|
.build()
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let req = build_request(input, Method::PUT);
|
||||||
|
let usecase = DefaultBucketUsecase::without_context();
|
||||||
|
|
||||||
|
let err = usecase.execute_put_bucket_notification_configuration(req).await.unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn execute_put_bucket_cors_returns_internal_error_when_store_uninitialized() {
|
async fn execute_put_bucket_cors_returns_internal_error_when_store_uninitialized() {
|
||||||
let input = PutBucketCorsInput::builder()
|
let input = PutBucketCorsInput::builder()
|
||||||
|
|||||||
Reference in New Issue
Block a user