Compare commits

..

3 Commits

Author SHA1 Message Date
Zhengchao An 707d062174 ci: stop running expensive jobs that cannot inform the result (#5528)
Three independent fixes that all avoid burning self-hosted runners on work
whose outcome is already determined. None of them changes what is tested.

- ci-docs-only: add a "Quick Checks" companion job. It is a prerequisite for
  gating ci.yml's expensive jobs behind quick-checks (rustfs/backlog#1599):
  once "Quick Checks" is a required check, a docs-only PR would otherwise wait
  on it forever. The steps are a byte-identical copy of ci.yml's quick-checks
  rather than an `echo`, so that on a mixed PR the two same-named check runs
  execute the same commands against the same merge ref and cannot disagree —
  GitHub has no written contract for how it picks between same-named required
  check runs, and the real job only takes 45-51s, leaving no timing margin to
  rely on.

- ci: guard uring-integration with the same `closed` check every other job
  already has. The pull_request trigger includes `closed` only so the
  concurrency group cancels in-flight runs; this job had no guard and no
  `needs`, so every closed or merged PR ran the full io_uring suite (4m17s,
  7m19s and 7m31s on runs 30678272341, 30678117601 and 30662728539).

- ci: gate s3-lifecycle-behavior-tests on e2e-tests, matching
  s3-implemented-tests. Both lanes only download the prebuilt debug binary, and
  s3-implemented-tests already finishes later, so a green PR's wall clock is
  unchanged; a red one stops holding a sm-standard-4 for up to 30 minutes.

Refs: rustfs/backlog#1598, rustfs/backlog#1599
2026-08-01 10:49:19 +08:00
houseme 4b6b6f14bd feat(hotpath): use mimalloc in inner counting allocator (#5523)
* feat(hotpath): use mimalloc in counting allocator

* fix(hotpath): adapt mimalloc for allocation counting

---------

Co-authored-by: overtrue <anzhengchao@gmail.com>
2026-08-01 10:43:52 +08:00
Zhengchao An 9080ea8ea0 test(object-lock): cover unretained version cleanup (#5504)
Co-authored-by: cxymds <cxymds@gmail.com>
2026-08-01 10:14:32 +08:00
7 changed files with 269 additions and 434 deletions
+63 -7
View File
@@ -12,18 +12,20 @@
# See the License for the specific language governing permissions and
# limitations under the License.
# Companion to ci.yml for the required "Test and Lint" status check.
# Companion to ci.yml for the required "Test and Lint" and "Quick Checks"
# status checks.
#
# ci.yml skips docs-only pull requests via paths-ignore, but the branch
# ruleset requires a check named "Test and Lint" — without this workflow a
# docs-only PR would wait on that check forever. This workflow triggers on
# exactly the paths ci.yml ignores and reports an instant success under the
# same job name. Mixed PRs trigger both workflows and the real check still
# gates: a required check with any failing run blocks the merge.
# ruleset requires checks named "Test and Lint" and "Quick Checks" — without
# this workflow a docs-only PR would wait on those checks forever. This
# workflow triggers on exactly the paths ci.yml ignores and reports success
# under the same job names. Mixed PRs trigger both workflows and the real
# checks still gate: a required check with any failing run blocks the merge.
# https://docs.github.com/en/repositories/configuring-branches-and-merges-in-your-repository/defining-the-mergeability-of-pull-requests/troubleshooting-required-status-checks#handling-skipped-but-required-checks
#
# Keep the paths list below in sync with the pull_request paths-ignore list
# in ci.yml.
# in ci.yml, and keep the quick-checks steps below byte-identical to the
# quick-checks job in ci.yml (see the comment on that job).
name: Continuous Integration (docs only)
@@ -52,9 +54,63 @@ permissions:
contents: read
jobs:
# Deliberately NOT a bare `echo`. Once "Quick Checks" becomes a required
# check, ci.yml gates every expensive job behind it, so a mixed PR reports
# two check runs with this name: the real one (45-51s) and this companion.
# GitHub has no written contract for how it picks between same-named
# required check runs ("latest wins" vs "any failure blocks"), so instead of
# relying on ordering we make both runs execute the same commands against
# the same merge ref — their conclusions are then necessarily identical and
# the choice does not matter. Keep these steps byte-identical to the
# quick-checks job in ci.yml (a guard script that asserts this, and the paths
# sync below, is tracked in rustfs/backlog#1603).
#
# For a genuinely docs-only PR this adds no strictness (no code changed, so
# fmt and the guards always pass) and costs ~50s of ubuntu-latest.
quick-checks:
name: Quick Checks
runs-on: ubuntu-latest
timeout-minutes: 10
steps:
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
- name: Install ripgrep
run: sudo apt-get update && sudo apt-get install -y ripgrep
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
with:
components: rustfmt
- name: Check code formatting
run: cargo fmt --all --check
- name: Check unsafe code allowances
run: ./scripts/check_unsafe_code_allowances.sh
- name: Check layered dependencies
run: ./scripts/check_layer_dependencies.sh
- name: Check architecture migration rules
run: ./scripts/check_architecture_migration_rules.sh
- name: Check tokio io-uring feature guard
run: ./scripts/check_no_tokio_io_uring.sh
- name: Check extension schema boundaries
run: ./scripts/check_extension_schema_boundaries.sh
- name: Check body-cache whitelist guard
run: ./scripts/check_body_cache_whitelist.sh
- name: Check no planning docs committed
run: ./scripts/check_no_planning_docs.sh
test-and-lint:
name: Test and Lint
runs-on: ubuntu-latest
timeout-minutes: 10
steps:
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
+14 -1
View File
@@ -448,6 +448,13 @@ jobs:
uring-integration:
name: io_uring Integration (real)
# The pull_request trigger includes `closed` purely so the concurrency
# group cancels in-flight runs of a closed PR; every other job opts out of
# that run with this guard (or is skipped through its `needs` chain). This
# job had neither, so each closed/merged PR really ran the whole io_uring
# suite (measured 4m17s / 7m19s / 7m31s on runs 30678272341 / 30678117601 /
# 30662728539) and kept the cancellation run in progress for minutes.
if: github.event_name != 'pull_request' || github.event.action != 'closed'
# GitHub-hosted ubuntu-latest runs a recent kernel with io_uring and, unlike
# a container, applies no seccomp filter that would block io_uring_setup — so
# the probe succeeds and the tests exercise the real UringBackend/FdCache/
@@ -746,7 +753,13 @@ jobs:
# evaluates ILM within ~2s of the due time, well inside the poll window.
s3-lifecycle-behavior-tests:
name: S3 Lifecycle Behavior Tests
needs: [ build-rustfs-debug-binary ]
# Also gated on e2e-tests, matching s3-implemented-tests: when the e2e smoke
# suite is already red this lane cannot tell us anything new, and it holds a
# sm-standard-4 for up to 30 minutes doing so. Both lanes only download the
# prebuilt debug binary (no cargo build), and s3-implemented-tests — which
# already waits on e2e-tests — finishes later anyway, so a green PR's total
# wall clock is unchanged.
needs: [ build-rustfs-debug-binary, e2e-tests ]
runs-on: sm-standard-4
timeout-minutes: 30
steps:
@@ -26,6 +26,7 @@
use super::common::*;
use aws_sdk_s3::Client;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::{ByteStream, DateTimeFormat};
use aws_sdk_s3::types::{
CompletedMultipartUpload, CompletedPart, Delete, MetadataDirective, ObjectIdentifier, ObjectLockLegalHoldStatus,
@@ -2120,6 +2121,127 @@ async fn test_multipart_default_retention_fixed_at_create() {
// Versioning Auto-Enable Tests
// ============================================================================
#[tokio::test]
#[serial]
async fn test_unretained_object_lock_object_delete_and_bucket_cleanup() {
init_logging();
info!("🧪 Test: Unretained Object Lock object delete and bucket cleanup (Issue #5339)");
let mut env = ObjectLockTestEnvironment::new()
.await
.expect("failed to create Object Lock test environment");
env.start_rustfs().await.expect("failed to start RustFS");
let bucket = "test-object-lock-delete-cleanup";
let key = "unretained-object";
env.create_object_lock_bucket(bucket)
.await
.expect("failed to create Object Lock bucket");
let client = env.s3_client();
let put_response = client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from_static(b"unretained data"))
.send()
.await
.expect("failed to upload unretained object");
let object_version_id = put_response
.version_id()
.expect("Object Lock buckets must create versioned objects")
.to_string();
let delete_response = client
.delete_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect("failed to create delete marker");
assert_eq!(delete_response.delete_marker(), Some(true));
let delete_marker_version_id = delete_response
.version_id()
.expect("Deleting without a version ID must create a delete marker")
.to_string();
let get_error = client
.get_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect_err("GET must not return an object hidden by a delete marker");
assert_eq!(get_error.raw_response().map(|response| response.status().as_u16()), Some(404));
assert_eq!(get_error.as_service_error().and_then(|error| error.code()), Some("NoSuchKey"));
let listed_objects = client
.list_objects_v2()
.bucket(bucket)
.send()
.await
.expect("failed to list current objects");
assert!(
listed_objects.contents().iter().all(|object| object.key() != Some(key)),
"ListObjectsV2 must hide objects whose latest version is a delete marker"
);
let listed_versions = client
.list_object_versions()
.bucket(bucket)
.send()
.await
.expect("failed to list object versions");
assert!(
listed_versions
.versions()
.iter()
.any(|version| version.key() == Some(key) && version.version_id() == Some(object_version_id.as_str())),
"The data version must remain until it is explicitly deleted"
);
assert!(
listed_versions
.delete_markers()
.iter()
.any(|marker| marker.key() == Some(key) && marker.version_id() == Some(delete_marker_version_id.as_str())),
"ListObjectVersions must expose the delete marker"
);
client
.delete_object()
.bucket(bucket)
.key(key)
.version_id(object_version_id)
.send()
.await
.expect("failed to delete the data version");
client
.delete_object()
.bucket(bucket)
.key(key)
.version_id(delete_marker_version_id)
.send()
.await
.expect("failed to delete the delete marker");
let remaining_versions = client
.list_object_versions()
.bucket(bucket)
.send()
.await
.expect("failed to list versions after cleanup");
assert!(remaining_versions.versions().is_empty());
assert!(remaining_versions.delete_markers().is_empty());
client
.delete_bucket()
.bucket(bucket)
.send()
.await
.expect("Deleting every version must remove xl.meta so the bucket can be deleted normally");
}
#[tokio::test]
#[serial]
async fn test_versioning_auto_enabled_with_object_lock() {
-54
View File
@@ -46,20 +46,6 @@ impl<'a> AdminResourceScope<'a> {
pub fn bucket_object(bucket: &'a str, object: &'a str) -> Self {
Self { bucket, object }
}
/// Scope an admin request to a single KMS key.
///
/// The policy crate carries the requested key identifier in the object slot
/// with an empty bucket (see `Statement::kms_key_scope_matches`); `Args` has
/// no dedicated field for it. An empty `key_id` means the caller could not
/// name a target key, which keeps the pre-resource-scoping behaviour where a
/// KMS statement matches every key.
pub fn kms_key(key_id: &'a str) -> Self {
Self {
bucket: "",
object: key_id,
}
}
}
pub async fn validate_admin_request(
@@ -188,34 +174,6 @@ pub async fn validate_admin_request_with_bucket_object(
evaluate_admin_actions(iam_store, &ctx, &actions, resource.bucket, resource.object).await
}
/// Admin gate for KMS endpoints that act on one key.
///
/// `key_id` is the identifier as requested (before any alias resolution), so a
/// policy scoped to `arn:aws:kms:::key/<id>` only authorizes that key. Endpoints
/// without a target key pass `""` and stay unscoped, which is also what a
/// malformed request resolves to: the request is rejected on its own parse error
/// right after the gate, so no key is ever touched under an unscoped decision.
pub async fn validate_admin_request_with_kms_key(
headers: &HeaderMap,
cred: &Credentials,
is_owner: bool,
deny_only: bool,
actions: Vec<Action>,
remote_addr: Option<std::net::SocketAddr>,
key_id: &str,
) -> S3Result<()> {
validate_admin_request_with_bucket_object(
headers,
cred,
is_owner,
deny_only,
actions,
remote_addr,
AdminResourceScope::kms_key(key_id),
)
.await
}
/// Unified authentication request handler for both UI and CLI
///
/// This function provides a single entry point for authentication,
@@ -513,18 +471,6 @@ mod tests {
assert_access_denied(res);
}
/// KMS scoping rides the object slot with an empty bucket, matching the
/// contract the policy crate evaluates KMS statements against.
#[test]
fn kms_scope_carries_the_key_id_in_the_object_slot() {
let scope = AdminResourceScope::kms_key("key-a");
assert_eq!(scope.bucket, "");
assert_eq!(scope.object, "key-a");
let unscoped = AdminResourceScope::kms_key("");
assert_eq!(unscoped.object, "", "an absent key id must stay unscoped");
}
/// The multi-action loop authorizes as soon as one candidate action passes
/// (owner short-circuits every action), and denies when none pass.
#[tokio::test]
+8 -56
View File
@@ -14,8 +14,8 @@
//! KMS key lifecycle admin API handlers: enable, disable and rotate.
use super::kms_keys::{extract_query_params, scoped_key_id};
use crate::admin::auth::validate_admin_request_with_kms_key;
use super::kms_keys::extract_query_params;
use crate::admin::auth::validate_admin_request;
use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::runtime_sources::current_kms_runtime_service_manager;
use crate::auth::{check_key_valid, get_session_token};
@@ -252,27 +252,21 @@ async fn handle_lifecycle_request(
let (cred, owner) = check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
// The body (or the `keyId` query parameter) names the key this operation acts
// on, so it has to be resolved before the gate runs. A read failure is
// surfaced only afterwards so an unauthorized caller still sees AccessDenied.
let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await;
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
actions,
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
body.as_ref()
.ok()
.and_then(|body| scoped_key_id(body, &req.uri))
.as_deref()
.unwrap_or_default(),
)
.await?;
let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let body = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let request: KmsKeyLifecycleRequest = if body.is_empty() {
let query_params = extract_query_params(&req.uri);
@@ -341,7 +335,6 @@ impl Operation for RotateKmsKeyHandler {
#[cfg(test)]
mod tests {
use super::*;
use http::Uri;
use rustfs_kms::backends::local::LocalKmsBackend;
use rustfs_kms::config::KmsConfig;
use rustfs_kms::types::{CreateKeyRequest, KeyState};
@@ -454,47 +447,6 @@ mod tests {
}
}
/// The lifecycle gate must be scoped to the key the request names, and the
/// key must be resolved the same way for authorization and for execution.
#[test]
fn lifecycle_requests_authorize_against_the_key_they_operate_on() {
let src = include_str!("kms_key_lifecycle.rs");
let handler = src
.split_once("async fn handle_lifecycle_request(")
.expect("lifecycle entry point should exist")
.1;
let handler = &handler[..handler.find("\nfn kms_service_manager_from_context").unwrap_or(handler.len())];
assert!(
handler.contains("validate_admin_request_with_kms_key("),
"lifecycle requests must scope authorization to the requested key"
);
assert!(
handler.contains("scoped_key_id(body, &req.uri)"),
"authorization must resolve the key exactly as execution does"
);
let uri: Uri = "/rustfs/admin/v3/kms/keys/disable?keyId=query-key"
.parse()
.expect("uri should parse");
let body = br#"{"key_id":"body-key"}"#;
assert_eq!(scoped_key_id(body, &uri).as_deref(), Some("body-key"));
assert_eq!(
scoped_key_id(body, &uri).as_deref(),
Some(
serde_json::from_slice::<KmsKeyLifecycleRequest>(body)
.expect("lifecycle body should parse")
.key_id
.as_str()
)
);
assert_eq!(
scoped_key_id(b"", &uri).as_deref(),
extract_query_params(&uri).get("keyId").map(String::as_str)
);
}
#[test]
fn lifecycle_request_rejects_unknown_fields() {
let error = serde_json::from_str::<KmsKeyLifecycleRequest>(r#"{"key_id":"key","unexpected_field":true}"#)
+26 -315
View File
@@ -14,7 +14,7 @@
//! KMS key management admin API handlers
use crate::admin::auth::{validate_admin_request, validate_admin_request_with_kms_key};
use crate::admin::auth::validate_admin_request;
use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::runtime_sources::{current_kms_runtime_service_manager, current_or_init_kms_runtime_service_manager};
use crate::auth::{check_key_valid, get_session_token};
@@ -116,29 +116,6 @@ fn extract_key_id(uri: &hyper::Uri) -> Option<String> {
.find_map(|name| query_params.get(name).filter(|value| !value.is_empty()).cloned())
}
/// The `key_id` of a KMS admin request body, read without committing to the
/// strict schema of the endpoint: the authorization gate needs the target key
/// before the body is parsed for execution, and a body that fails the strict
/// parse is rejected right after the gate anyway.
#[derive(Deserialize)]
struct KeyIdProbe {
#[serde(default)]
key_id: String,
}
/// Target key of an endpoint that accepts either a JSON body or a `keyId` query
/// parameter, resolved exactly the way the endpoint resolves it for execution so
/// the authorized key can never differ from the operated one. Returns `None`
/// when the target cannot be determined, which authorizes unscoped and lets the
/// request fail on its own parse error afterwards.
pub(super) fn scoped_key_id(body: &[u8], uri: &hyper::Uri) -> Option<String> {
if body.is_empty() {
return extract_query_params(uri).get("keyId").cloned();
}
serde_json::from_slice::<KeyIdProbe>(body).ok().map(|probe| probe.key_id)
}
fn kms_service_manager_from_context() -> Option<std::sync::Arc<rustfs_kms::KmsServiceManager>> {
current_kms_runtime_service_manager()
}
@@ -302,20 +279,17 @@ impl Operation for DescribeKeyHandler {
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
let requested_key_id = extract_key_id(&req.uri);
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
kms_describe_key_actions(),
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
requested_key_id.as_deref().unwrap_or_default(),
)
.await?;
let Some(key_id) = requested_key_id else {
let Some(key_id) = extract_key_id(&req.uri) else {
return Err(s3_error!(InvalidRequest, "missing required parameter: 'keyId'"));
};
@@ -361,12 +335,10 @@ mod tests {
use super::{
CancelKmsKeyDeletionRequest, CreateKeyApiRequest, CreateKmsKeyRequest, DeleteKmsKeyRequest, GenerateDataKeyApiRequest,
extract_key_id, kms_create_key_actions, kms_delete_key_actions, kms_describe_key_actions, kms_generate_data_key_actions,
kms_list_keys_actions, scoped_key_id,
kms_list_keys_actions,
};
use http::Uri;
use rustfs_policy::policy::action::{Action, AdminAction, KmsAction};
use rustfs_policy::policy::{Args, Policy};
use std::collections::HashMap;
fn assert_has_action(actions: &[Action], action: Action) {
assert!(actions.contains(&action), "expected action list to contain {action:?}");
@@ -450,248 +422,6 @@ mod tests {
assert!(err.to_string().contains("unknown field"));
}
}
#[test]
fn scoped_key_id_reads_the_body_first_and_falls_back_to_the_query() {
let uri: Uri = "/rustfs/admin/v3/kms/keys/delete?keyId=query-key"
.parse()
.expect("uri should parse");
assert_eq!(scoped_key_id(br#"{"key_id":"body-key"}"#, &uri).as_deref(), Some("body-key"));
assert_eq!(scoped_key_id(b"", &uri).as_deref(), Some("query-key"));
// The query aliases accepted by the legacy describe endpoint are not
// accepted here, because the handlers only execute on `keyId`.
let alias_uri: Uri = "/rustfs/admin/v3/kms/keys/delete?key-id=query-key"
.parse()
.expect("uri should parse");
assert_eq!(scoped_key_id(b"", &alias_uri), None);
}
/// A body the endpoint will reject leaves the request unscoped: the gate then
/// answers as it did before resource scoping and the request still fails on
/// its own parse error, so no key is operated on under that decision.
#[test]
fn scoped_key_id_is_absent_for_bodies_that_name_no_key() {
let uri: Uri = "/rustfs/admin/v3/kms/keys/delete".parse().expect("uri should parse");
assert_eq!(scoped_key_id(b"not json", &uri), None);
assert_eq!(scoped_key_id(br#"{"key_id":"" }"#, &uri).as_deref(), Some(""));
assert_eq!(scoped_key_id(b"{}", &uri).as_deref(), Some(""));
assert_eq!(scoped_key_id(b"", &uri), None);
}
/// Identity policy scoped to a single KMS key, or unscoped when `key_arn` is
/// `None` (the legacy action-only form that must keep matching every key).
fn kms_policy(key_arn: Option<&str>) -> Policy {
let resource = match key_arn {
Some(arn) => format!(r#","Resource":["{arn}"]"#),
None => String::new(),
};
let document = format!(r#"{{"Version":"2012-10-17","Statement":[{{"Effect":"Allow","Action":["kms:*"]{resource}}}]}}"#);
Policy::parse_config(document.as_bytes()).expect("kms policy should parse")
}
async fn policy_allows(policy: &Policy, action: Action, key_id: &str, is_owner: bool) -> bool {
let groups = None;
let conditions = HashMap::new();
let claims = HashMap::new();
policy
.is_allowed(&Args {
account: "kms-operator",
groups: &groups,
action,
// The admin gate scopes KMS requests through `AdminResourceScope::kms_key`,
// which puts the requested key id in the object slot with an empty bucket.
bucket: "",
conditions: &conditions,
is_owner,
object: key_id,
claims: &claims,
deny_only: false,
})
.await
}
fn describe_by_query(key_id: &str) -> String {
let uri: Uri = format!("/rustfs/admin/v3/kms/key/status?keyId={key_id}")
.parse()
.expect("uri should parse");
extract_key_id(&uri).unwrap_or_default()
}
fn describe_by_path(key_id: &str) -> String {
// `DescribeKmsKeyHandler` reads the router path parameter verbatim.
key_id.to_string()
}
fn generate_data_key_body(key_id: &str) -> String {
let body = format!(r#"{{"key_id":"{key_id}","key_spec":"Aes256"}}"#);
serde_json::from_str::<GenerateDataKeyApiRequest>(&body)
.expect("generate-data-key body should parse")
.key_id
}
fn key_id_body(key_id: &str) -> String {
let uri: Uri = "/rustfs/admin/v3/kms/keys".parse().expect("uri should parse");
scoped_key_id(format!(r#"{{"key_id":"{key_id}"}}"#).as_bytes(), &uri).unwrap_or_default()
}
fn key_id_query(key_id: &str) -> String {
let uri: Uri = format!("/rustfs/admin/v3/kms/keys?keyId={key_id}")
.parse()
.expect("uri should parse");
scoped_key_id(b"", &uri).unwrap_or_default()
}
/// Endpoint label, the action it gates on, and the resolution its handler
/// performs to obtain the key the gate is scoped to.
type SingleKeyEndpoint = (&'static str, Action, fn(&str) -> String);
/// Every KMS admin endpoint that acts on one key. Lifecycle endpoints share
/// `scoped_key_id` with the delete/cancel endpoints (see `kms_key_lifecycle`).
fn single_key_endpoints() -> Vec<SingleKeyEndpoint> {
vec![
(
"GET /v3/kms/key/status",
Action::KmsAction(KmsAction::DescribeKeyAction),
describe_by_query as fn(&str) -> String,
),
(
"GET /v3/kms/keys/{key_id}",
Action::KmsAction(KmsAction::DescribeKeyAction),
describe_by_path,
),
(
"POST /v3/kms/generate-data-key",
Action::KmsAction(KmsAction::GenerateDataKeyAction),
generate_data_key_body,
),
("DELETE /v3/kms/keys/delete", Action::KmsAction(KmsAction::DeleteKeyAction), key_id_body),
(
"DELETE /v3/kms/keys/delete?keyId=",
Action::KmsAction(KmsAction::DeleteKeyAction),
key_id_query,
),
(
"POST /v3/kms/keys/cancel-deletion",
Action::KmsAction(KmsAction::DeleteKeyAction),
key_id_body,
),
("POST /v3/kms/keys/enable", Action::KmsAction(KmsAction::EnableKeyAction), key_id_body),
("POST /v3/kms/keys/disable", Action::KmsAction(KmsAction::DisableKeyAction), key_id_body),
("POST /v3/kms/keys/rotate", Action::KmsAction(KmsAction::RotateKeyAction), key_id_body),
]
}
/// A policy limited to `key/A` must not authorize any single-key endpoint
/// against another key (rustfs/backlog#1582).
#[tokio::test]
async fn single_key_endpoints_reject_a_key_outside_the_policy_scope() {
let policy = kms_policy(Some("arn:aws:kms:::key/key-a"));
for (endpoint, action, resolve_key_id) in single_key_endpoints() {
assert!(
policy_allows(&policy, action, &resolve_key_id("key-a"), false).await,
"{endpoint} must stay allowed on the key the policy names"
);
assert!(
!policy_allows(&policy, action, &resolve_key_id("key-b"), false).await,
"{endpoint} must be denied on a key outside the policy scope"
);
}
}
/// Compatibility pin: a KMS statement without resources keeps matching every
/// key, and the owner short-circuit is unaffected by scoping.
#[tokio::test]
async fn unscoped_policies_and_owner_credentials_keep_matching_every_key() {
let unscoped = kms_policy(None);
let scoped = kms_policy(Some("arn:aws:kms:::key/key-a"));
for (endpoint, action, resolve_key_id) in single_key_endpoints() {
for key in ["key-a", "key-b"] {
assert!(
policy_allows(&unscoped, action, &resolve_key_id(key), false).await,
"{endpoint} must stay allowed on {key} for a resource-less KMS statement"
);
assert!(
policy_allows(&scoped, action, &resolve_key_id(key), true).await,
"{endpoint} must stay allowed on {key} for an owner credential"
);
}
}
}
/// The key the handler resolves also reaches Deny statements, so an explicit
/// Deny on one key survives an unscoped Allow.
#[tokio::test]
async fn explicit_deny_on_one_key_overrides_an_unscoped_allow() {
let document = r#"{"Version":"2012-10-17","Statement":[
{"Effect":"Allow","Action":["kms:*"]},
{"Effect":"Deny","Action":["kms:*"],"Resource":["arn:aws:kms:::key/key-b"]}
]}"#;
let policy = Policy::parse_config(document.as_bytes()).expect("kms policy should parse");
for (endpoint, action, resolve_key_id) in single_key_endpoints() {
assert!(
policy_allows(&policy, action, &resolve_key_id("key-a"), false).await,
"{endpoint} must stay allowed on a key the Deny does not name"
);
assert!(
!policy_allows(&policy, action, &resolve_key_id("key-b"), false).await,
"{endpoint} must be denied on the key the Deny names"
);
}
}
/// The single-key handlers must reach the gate through the KMS-scoped entry
/// point; the endpoints that have no target key must not (passing an empty
/// scope there would be a no-op, but the unscoped call documents intent).
#[test]
fn single_key_handlers_authorize_through_the_kms_scoped_gate() {
let src = include_str!("kms_keys.rs");
for handler in [
"DescribeKeyHandler",
"GenerateDataKeyHandler",
"DeleteKmsKeyHandler",
"CancelKmsKeyDeletionHandler",
"DescribeKmsKeyHandler",
] {
let block = operation_block(src, handler);
assert!(
block.contains("validate_admin_request_with_kms_key("),
"{handler} must scope its authorization to the requested key"
);
}
for handler in [
"CreateKeyHandler",
"ListKeysHandler",
"CreateKmsKeyHandler",
"ListKmsKeysHandler",
] {
let block = operation_block(src, handler);
assert!(
!block.contains("validate_admin_request_with_kms_key("),
"{handler} has no target key and must not claim a KMS resource scope"
);
}
}
fn operation_block<'a>(src: &'a str, handler: &str) -> &'a str {
let marker = format!("impl Operation for {handler}");
let block = src.split_once(&marker).expect("handler impl should exist").1;
let end = ["\n/// ", "\n#[derive(", "\n#[cfg(test)]", "\npub struct "]
.into_iter()
.filter_map(|item| block.find(item))
.min()
.unwrap_or(block.len());
&block[..end]
}
}
/// List KMS keys (legacy endpoint)
@@ -777,31 +507,24 @@ impl Operation for GenerateDataKeyHandler {
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
// The body names the key this endpoint derives a data key from, so it has
// to be read before the gate runs. Input failures stay deferred until
// after the gate so an unauthorized caller still sees AccessDenied.
let parsed = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))
.and_then(|body| {
serde_json::from_slice::<GenerateDataKeyApiRequest>(&body)
.map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))
});
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
kms_generate_data_key_actions(),
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
parsed.as_ref().map(|request| request.key_id.as_str()).unwrap_or_default(),
)
.await?;
let request = parsed?;
let body = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let request: GenerateDataKeyApiRequest =
serde_json::from_slice(&body).map_err(|e| s3_error!(InvalidRequest, "invalid JSON: {}", e))?;
let Some(service) = kms_encryption_service_from_context().await else {
return Err(s3_error!(InternalError, "KMS service not initialized"));
@@ -1009,27 +732,21 @@ impl Operation for DeleteKmsKeyHandler {
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
// Read the body before the gate so the request can be scoped to the key it
// targets; the read failure is surfaced afterwards to keep AccessDenied
// ahead of input errors for callers that are not authorized at all.
let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await;
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
kms_delete_key_actions(),
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
body.as_ref()
.ok()
.and_then(|body| scoped_key_id(body, &req.uri))
.as_deref()
.unwrap_or_default(),
)
.await?;
let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let body = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let request: DeleteKmsKeyRequest = if body.is_empty() {
let query_params = extract_query_params(&req.uri);
@@ -1189,26 +906,21 @@ impl Operation for CancelKmsKeyDeletionHandler {
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
// Same ordering as the delete endpoint: the target key is resolved before
// the gate, the read failure only after it.
let body = req.input.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE).await;
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
kms_delete_key_actions(),
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
body.as_ref()
.ok()
.and_then(|body| scoped_key_id(body, &req.uri))
.as_deref()
.unwrap_or_default(),
)
.await?;
let body = body.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let body = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|e| s3_error!(InvalidRequest, "failed to read request body: {}", e))?;
let request: CancelKmsKeyDeletionRequest = if body.is_empty() {
let query_params = extract_query_params(&req.uri);
@@ -1468,14 +1180,13 @@ impl Operation for DescribeKmsKeyHandler {
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
validate_admin_request_with_kms_key(
validate_admin_request(
&req.headers,
&cred,
owner,
false,
kms_describe_key_actions(),
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
params.get("key_id").unwrap_or_default(),
)
.await?;
+36 -1
View File
@@ -12,9 +12,32 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
use std::alloc::{GlobalAlloc, Layout};
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
#[derive(Default)]
struct DefaultMiMalloc;
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
// SAFETY: allocation and deallocation are forwarded unchanged to MiMalloc, so
// MiMalloc's GlobalAlloc guarantees apply to every returned pointer and layout.
#[allow(unsafe_code)]
unsafe impl GlobalAlloc for DefaultMiMalloc {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { mimalloc::MiMalloc.alloc(layout) }
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
}
}
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
#[global_allocator]
static GLOBAL: hotpath::CountingAllocator = hotpath::CountingAllocator::new();
static GLOBAL: hotpath::CountingAllocator<DefaultMiMalloc> = hotpath::CountingAllocator::new();
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
#[global_allocator]
@@ -25,3 +48,15 @@ fn main() {
rustfs::startup_entrypoint::run_process();
}
#[cfg(all(test, feature = "hotpath", feature = "hotpath-alloc"))]
mod tests {
#[test]
#[allow(unsafe_code)]
fn hotpath_allocator_uses_mimalloc() {
let allocation = Box::new([0_u8; 64]);
// SAFETY: the live Box pointer is valid to inspect for heap ownership.
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
}
}