mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
feat(ilm): add manual transition run endpoint (#5171)
* feat(ilm): report manual transition backfill outcomes Add scoped lifecycle transition backfill reporting for backlog #1478 and expose enqueue outcomes needed by #1479 without changing the existing scanner/compensation bool API. Co-Authored-By: heihutu <heihutu@gmail.com> * feat(admin): add manual transition run endpoint Add a bounded POST /rustfs/admin/v3/ilm/transition/run API for backlog #1477 and cover the route, policy, query parsing, and partial status contract needed by #1481. Console operations from #1480 are intentionally left for a later client integration. Co-Authored-By: heihutu <heihutu@gmail.com> * feat(ilm): allow manual transition resume markers Accept additive marker and versionMarker parameters on the bounded manual transition run API so clients can continue from a partial report without changing the existing default scan behavior. Co-Authored-By: heihutu <heihutu@gmail.com> * fix(ilm): harden manual transition partial reports Preserve the null-version cursor contract, stop manual scans on enqueue pressure without skipping the failed object, and keep raw resume markers out of admin JSON responses. Co-Authored-By: heihutu <heihutu@gmail.com> * test(ilm): add manual transition e2e coverage Co-Authored-By: heihutu <heihutu@gmail.com> * perf(ilm): keep transition enqueue hot path direct Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -48,15 +48,17 @@ use aws_sdk_s3::Client;
|
||||
use aws_sdk_s3::error::ProvideErrorMetadata;
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use aws_sdk_s3::types::{
|
||||
BucketLifecycleConfiguration, CompletedMultipartUpload, CompletedPart, ExpirationStatus, LifecycleRule, LifecycleRuleFilter,
|
||||
Transition, TransitionStorageClass,
|
||||
BucketLifecycleConfiguration, BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, ExpirationStatus,
|
||||
LifecycleRule, LifecycleRuleFilter, NoncurrentVersionTransition, Transition, TransitionStorageClass, VersioningConfiguration,
|
||||
};
|
||||
use http::Method;
|
||||
use http::header::HOST;
|
||||
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
|
||||
use rustfs_signer::sign_v4;
|
||||
use s3s::Body;
|
||||
use serde::Deserialize;
|
||||
use std::time::{Duration as StdDuration, Instant};
|
||||
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
|
||||
|
||||
type TestResult = Result<(), Box<dyn std::error::Error + Send + Sync>>;
|
||||
|
||||
@@ -64,10 +66,18 @@ const TIER_NAME: &str = "COLDTIER";
|
||||
const TIER_BUCKET: &str = "ilm7-cold-tier";
|
||||
const TIER_PREFIX: &str = "tiered";
|
||||
const SOURCE_BUCKET: &str = "ilm7-hot";
|
||||
const MANUAL_DUE_BUCKET: &str = "ilm7-manual-due";
|
||||
const MANUAL_DRY_RUN_BUCKET: &str = "ilm7-manual-dry-run";
|
||||
const MANUAL_NOT_DUE_BUCKET: &str = "ilm7-manual-not-due";
|
||||
const OBJECT_KEY: &str = "tier/鲁A12345/report.bin";
|
||||
const MANUAL_DUE_KEY: &str = "manual-due/report.bin";
|
||||
const MANUAL_DRY_RUN_KEY: &str = "manual-dry-run/report.bin";
|
||||
const MANUAL_NOT_DUE_KEY: &str = "manual-not-due/report.bin";
|
||||
const CONTENT_TYPE: &str = "application/x-ilm7";
|
||||
const USER_META_KEY: &str = "ilm7-origin";
|
||||
const USER_META_VAL: &str = "hermetic-transition";
|
||||
const HDR_SOURCE_REPLICATION_REQUEST: &str = "x-rustfs-source-replication-request";
|
||||
const HDR_SOURCE_MTIME: &str = "x-rustfs-source-mtime";
|
||||
|
||||
/// 5 MiB — the S3 minimum size for a non-final multipart part; the object's only
|
||||
/// internal part boundary sits at this offset.
|
||||
@@ -158,12 +168,16 @@ async fn add_rustfs_tier(hot: &RustFSTestEnvironment, cold: &RustFSTestEnvironme
|
||||
|
||||
/// A current-version `Transition Days=0` rule scoped to the object's prefix.
|
||||
fn transition_rule() -> Result<LifecycleRule, Box<dyn std::error::Error + Send + Sync>> {
|
||||
transition_rule_for("ilm7-transition", "tier/", 0)
|
||||
}
|
||||
|
||||
fn transition_rule_for(id: &str, prefix: &str, days: i32) -> Result<LifecycleRule, Box<dyn std::error::Error + Send + Sync>> {
|
||||
Ok(LifecycleRule::builder()
|
||||
.id("ilm7-transition")
|
||||
.filter(LifecycleRuleFilter::builder().prefix("tier/").build())
|
||||
.id(id)
|
||||
.filter(LifecycleRuleFilter::builder().prefix(prefix).build())
|
||||
.transitions(
|
||||
Transition::builder()
|
||||
.days(0)
|
||||
.days(days)
|
||||
.storage_class(TransitionStorageClass::from(TIER_NAME))
|
||||
.build(),
|
||||
)
|
||||
@@ -171,6 +185,55 @@ fn transition_rule() -> Result<LifecycleRule, Box<dyn std::error::Error + Send +
|
||||
.build()?)
|
||||
}
|
||||
|
||||
async fn put_lifecycle_transition_rule(client: &Client, bucket: &str, id: &str, prefix: &str, days: i32) -> TestResult {
|
||||
let lifecycle = BucketLifecycleConfiguration::builder()
|
||||
.rules(transition_rule_for(id, prefix, days)?)
|
||||
.build()?;
|
||||
client
|
||||
.put_bucket_lifecycle_configuration()
|
||||
.bucket(bucket)
|
||||
.lifecycle_configuration(lifecycle)
|
||||
.send()
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn put_lifecycle_noncurrent_transition_rule(client: &Client, bucket: &str, id: &str, prefix: &str) -> TestResult {
|
||||
let rule = LifecycleRule::builder()
|
||||
.id(id)
|
||||
.filter(LifecycleRuleFilter::builder().prefix(prefix).build())
|
||||
.noncurrent_version_transitions(
|
||||
NoncurrentVersionTransition::builder()
|
||||
.noncurrent_days(0)
|
||||
.storage_class(TransitionStorageClass::from(TIER_NAME))
|
||||
.build(),
|
||||
)
|
||||
.status(ExpirationStatus::Enabled)
|
||||
.build()?;
|
||||
let lifecycle = BucketLifecycleConfiguration::builder().rules(rule).build()?;
|
||||
client
|
||||
.put_bucket_lifecycle_configuration()
|
||||
.bucket(bucket)
|
||||
.lifecycle_configuration(lifecycle)
|
||||
.send()
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn enable_bucket_versioning(client: &Client, bucket: &str) -> TestResult {
|
||||
client
|
||||
.put_bucket_versioning()
|
||||
.bucket(bucket)
|
||||
.versioning_configuration(
|
||||
VersioningConfiguration::builder()
|
||||
.status(BucketVersioningStatus::Enabled)
|
||||
.build(),
|
||||
)
|
||||
.send()
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Upload `data` as a two-part multipart object with a content-type and one
|
||||
/// user-metadata entry.
|
||||
async fn put_multipart_object(client: &Client, bucket: &str, key: &str, data: &[u8]) -> TestResult {
|
||||
@@ -218,6 +281,90 @@ async fn put_multipart_object(client: &Client, bucket: &str, key: &str, data: &[
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn put_single_part_object(client: &Client, bucket: &str, key: &str, body: &'static [u8]) -> TestResult {
|
||||
client
|
||||
.put_object()
|
||||
.bucket(bucket)
|
||||
.key(key)
|
||||
.body(ByteStream::from_static(body))
|
||||
.send()
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn put_backdated_single_part_object(
|
||||
client: &Client,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
body: &'static [u8],
|
||||
mtime: OffsetDateTime,
|
||||
) -> TestResult {
|
||||
let mtime_rfc3339 = mtime.format(&Rfc3339)?;
|
||||
client
|
||||
.put_object()
|
||||
.bucket(bucket)
|
||||
.key(key)
|
||||
.body(ByteStream::from_static(body))
|
||||
.customize()
|
||||
.mutate_request(move |req| {
|
||||
req.headers_mut().insert(HDR_SOURCE_REPLICATION_REQUEST, "true");
|
||||
req.headers_mut().insert(HDR_SOURCE_MTIME, mtime_rfc3339.clone());
|
||||
})
|
||||
.send()
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct ManualTransitionRunResponse {
|
||||
state: String,
|
||||
mode: String,
|
||||
job_id: Option<String>,
|
||||
status_endpoint: Option<String>,
|
||||
report: ManualTransitionRunReport,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct ManualTransitionRunReport {
|
||||
bucket: String,
|
||||
prefix: String,
|
||||
tier: Option<String>,
|
||||
dry_run: bool,
|
||||
lifecycle_config_found: bool,
|
||||
scanned: u64,
|
||||
eligible: u64,
|
||||
enqueued: u64,
|
||||
dry_run_eligible: u64,
|
||||
skipped_not_transition: u64,
|
||||
skipped_tier: u64,
|
||||
skipped_delete_marker: u64,
|
||||
skipped_directory: u64,
|
||||
skipped_replication: u64,
|
||||
skipped_already_in_flight: u64,
|
||||
skipped_queue_full: u64,
|
||||
skipped_queue_closed: u64,
|
||||
skipped_queue_timeout: u64,
|
||||
truncated_by_limit: bool,
|
||||
}
|
||||
|
||||
async fn manual_transition_run(
|
||||
hot: &RustFSTestEnvironment,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
dry_run: bool,
|
||||
) -> Result<ManualTransitionRunResponse, Box<dyn std::error::Error + Send + Sync>> {
|
||||
let bucket = urlencoding::encode(bucket);
|
||||
let prefix = urlencoding::encode(prefix);
|
||||
let tier = urlencoding::encode(TIER_NAME);
|
||||
let path =
|
||||
format!("/rustfs/admin/v3/ilm/transition/run?bucket={bucket}&prefix={prefix}&tier={tier}&dryRun={dry_run}&maxObjects=10");
|
||||
let (status, body) = signed_admin_request(&hot.url, Method::POST, &path, None, &hot.access_key, &hot.secret_key).await?;
|
||||
if !status.is_success() {
|
||||
return Err(format!("manual transition run failed: status={status}, body={body}").into());
|
||||
}
|
||||
Ok(serde_json::from_str(&body)?)
|
||||
}
|
||||
|
||||
/// Number of objects currently stored in the cold-tier bucket.
|
||||
async fn cold_tier_object_count(cold_client: &Client) -> Result<usize, Box<dyn std::error::Error + Send + Sync>> {
|
||||
let resp = cold_client.list_objects_v2().bucket(TIER_BUCKET).send().await?;
|
||||
@@ -283,6 +430,27 @@ async fn assert_range(client: &Client, start: usize, end: usize, data: &[u8]) ->
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn assert_not_transitioned(client: &Client, bucket: &str, key: &str) -> TestResult {
|
||||
let head = client.head_object().bucket(bucket).key(key).send().await?;
|
||||
assert!(
|
||||
head.storage_class().is_none(),
|
||||
"{bucket}/{key} must remain in the hot tier, got storage_class={:?}",
|
||||
head.storage_class()
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn assert_remains_not_transitioned(client: &Client, bucket: &str, key: &str, duration: StdDuration) -> TestResult {
|
||||
let deadline = Instant::now() + duration;
|
||||
loop {
|
||||
assert_not_transitioned(client, bucket, key).await?;
|
||||
if Instant::now() >= deadline {
|
||||
return Ok(());
|
||||
}
|
||||
tokio::time::sleep(StdDuration::from_millis(250)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Full ilm-7 hermetic transition main path across two embedded RustFS servers.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn test_hermetic_transition_main_path() -> TestResult {
|
||||
@@ -378,3 +546,100 @@ async fn test_hermetic_transition_main_path() -> TestResult {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn test_manual_transition_run_black_box_semantics() -> TestResult {
|
||||
let mut cold = RustFSTestEnvironment::new().await?;
|
||||
cold.access_key = "manualcoldtieradmin".to_string();
|
||||
cold.secret_key = "manualcoldtiersecret".to_string();
|
||||
cold.start_rustfs_server_without_cleanup(vec![]).await?;
|
||||
let cold_client = cold.create_s3_client();
|
||||
cold_client.create_bucket().bucket(TIER_BUCKET).send().await?;
|
||||
|
||||
let mut hot = RustFSTestEnvironment::new().await?;
|
||||
hot.start_rustfs_server_with_env(vec![], &[("RUSTFS_SCANNER_ENABLED", "false"), ("RUSTFS_SCANNER_CYCLE", "3600")])
|
||||
.await?;
|
||||
let hot_client = hot.create_s3_client();
|
||||
add_rustfs_tier(&hot, &cold).await?;
|
||||
let due_mtime = OffsetDateTime::now_utc() - time::Duration::hours(25);
|
||||
|
||||
hot_client.create_bucket().bucket(MANUAL_DUE_BUCKET).send().await?;
|
||||
put_lifecycle_transition_rule(&hot_client, MANUAL_DUE_BUCKET, "manual-due", "manual-due/", 0).await?;
|
||||
put_backdated_single_part_object(&hot_client, MANUAL_DUE_BUCKET, MANUAL_DUE_KEY, b"manual due object", due_mtime).await?;
|
||||
|
||||
let due = manual_transition_run(&hot, MANUAL_DUE_BUCKET, "manual-due/", false).await?;
|
||||
assert_eq!(due.mode, "enqueue_only");
|
||||
assert!(due.job_id.is_none());
|
||||
assert!(due.status_endpoint.is_none());
|
||||
assert_eq!(due.state, "completed");
|
||||
assert_eq!(due.report.bucket, MANUAL_DUE_BUCKET);
|
||||
assert_eq!(due.report.prefix, "manual-due/");
|
||||
assert_eq!(due.report.tier.as_deref(), Some(TIER_NAME));
|
||||
assert!(!due.report.dry_run);
|
||||
assert!(due.report.lifecycle_config_found);
|
||||
assert_eq!(due.report.scanned, 1, "due report: {:#?}", due.report);
|
||||
assert_eq!(due.report.eligible, 1, "due report: {:#?}", due.report);
|
||||
assert_eq!(
|
||||
due.report.enqueued + due.report.skipped_already_in_flight,
|
||||
1,
|
||||
"due report: {:#?}",
|
||||
due.report
|
||||
);
|
||||
assert_eq!(due.report.skipped_tier, 0);
|
||||
assert_eq!(due.report.skipped_delete_marker, 0);
|
||||
assert_eq!(due.report.skipped_directory, 0);
|
||||
assert_eq!(due.report.skipped_replication, 0);
|
||||
assert!(!due.report.truncated_by_limit);
|
||||
wait_for_transition(&hot_client, MANUAL_DUE_BUCKET, MANUAL_DUE_KEY, StdDuration::from_secs(90)).await?;
|
||||
let remote_count_after_due = cold_tier_object_count(&cold_client).await?;
|
||||
|
||||
hot_client.create_bucket().bucket(MANUAL_DRY_RUN_BUCKET).send().await?;
|
||||
enable_bucket_versioning(&hot_client, MANUAL_DRY_RUN_BUCKET).await?;
|
||||
put_lifecycle_noncurrent_transition_rule(&hot_client, MANUAL_DRY_RUN_BUCKET, "manual-dry-run", "manual-dry-run/").await?;
|
||||
put_single_part_object(&hot_client, MANUAL_DRY_RUN_BUCKET, MANUAL_DRY_RUN_KEY, b"manual dry-run object v1").await?;
|
||||
put_single_part_object(&hot_client, MANUAL_DRY_RUN_BUCKET, MANUAL_DRY_RUN_KEY, b"manual dry-run object v2").await?;
|
||||
|
||||
let before_dry_run_remote_count = cold_tier_object_count(&cold_client).await?;
|
||||
assert_eq!(
|
||||
before_dry_run_remote_count, remote_count_after_due,
|
||||
"dry-run setup must not enqueue transition work before the manual dry-run"
|
||||
);
|
||||
let dry = manual_transition_run(&hot, MANUAL_DRY_RUN_BUCKET, "manual-dry-run/", true).await?;
|
||||
assert_eq!(dry.state, "completed");
|
||||
assert_eq!(dry.report.bucket, MANUAL_DRY_RUN_BUCKET);
|
||||
assert_eq!(dry.report.prefix, "manual-dry-run/");
|
||||
assert_eq!(dry.report.tier.as_deref(), Some(TIER_NAME));
|
||||
assert!(dry.report.dry_run);
|
||||
assert_eq!(dry.report.scanned, 2, "dry-run report: {:#?}", dry.report);
|
||||
assert_eq!(dry.report.eligible, 1, "dry-run report: {:#?}", dry.report);
|
||||
assert_eq!(dry.report.dry_run_eligible, 1, "dry-run report: {:#?}", dry.report);
|
||||
assert_eq!(dry.report.enqueued, 0, "dry-run report: {:#?}", dry.report);
|
||||
assert_eq!(dry.report.skipped_not_transition, 1, "dry-run report: {:#?}", dry.report);
|
||||
assert_eq!(
|
||||
cold_tier_object_count(&cold_client).await?,
|
||||
before_dry_run_remote_count,
|
||||
"dry-run must not create a remote tier object"
|
||||
);
|
||||
assert_not_transitioned(&hot_client, MANUAL_DRY_RUN_BUCKET, MANUAL_DRY_RUN_KEY).await?;
|
||||
|
||||
hot_client.create_bucket().bucket(MANUAL_NOT_DUE_BUCKET).send().await?;
|
||||
put_lifecycle_transition_rule(&hot_client, MANUAL_NOT_DUE_BUCKET, "manual-not-due", "manual-not-due/", 1).await?;
|
||||
put_single_part_object(&hot_client, MANUAL_NOT_DUE_BUCKET, MANUAL_NOT_DUE_KEY, b"manual not-yet-due object").await?;
|
||||
|
||||
let not_due = manual_transition_run(&hot, MANUAL_NOT_DUE_BUCKET, "manual-not-due/", false).await?;
|
||||
assert_eq!(not_due.state, "completed");
|
||||
assert_eq!(not_due.report.bucket, MANUAL_NOT_DUE_BUCKET);
|
||||
assert_eq!(not_due.report.prefix, "manual-not-due/");
|
||||
assert_eq!(not_due.report.tier.as_deref(), Some(TIER_NAME));
|
||||
assert!(!not_due.report.dry_run);
|
||||
assert_eq!(not_due.report.scanned, 1, "not-due report: {:#?}", not_due.report);
|
||||
assert_eq!(not_due.report.eligible, 0, "not-due report: {:#?}", not_due.report);
|
||||
assert_eq!(not_due.report.enqueued, 0, "not-due report: {:#?}", not_due.report);
|
||||
assert_eq!(not_due.report.skipped_not_transition, 1, "not-due report: {:#?}", not_due.report);
|
||||
assert_eq!(not_due.report.skipped_queue_full, 0);
|
||||
assert_eq!(not_due.report.skipped_queue_closed, 0);
|
||||
assert_eq!(not_due.report.skipped_queue_timeout, 0);
|
||||
assert_remains_not_transitioned(&hot_client, MANUAL_NOT_DUE_BUCKET, MANUAL_NOT_DUE_KEY, StdDuration::from_secs(2)).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -43,10 +43,12 @@ pub mod bucket {
|
||||
|
||||
pub mod bucket_lifecycle_ops {
|
||||
pub use crate::bucket::lifecycle::bucket_lifecycle_ops::{
|
||||
ExpiryState, LifecycleOps, RestoreRequestOps, TransitionState, TransitionedObject, apply_expiry_rule,
|
||||
apply_transition_rule, enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects,
|
||||
enqueue_transition_immediate, expire_transitioned_object, get_global_expiry_state, get_global_transition_state,
|
||||
init_background_expiry, post_restore_opts, run_stale_multipart_upload_cleanup_once, validate_transition_tier,
|
||||
ExpiryState, LifecycleOps, ManualTransitionRunOptions, ManualTransitionRunReport, RestoreRequestOps,
|
||||
TransitionState, TransitionedObject, apply_expiry_rule, apply_transition_rule,
|
||||
enqueue_expiry_for_existing_objects, enqueue_transition_for_existing_objects,
|
||||
enqueue_transition_for_existing_objects_scoped, enqueue_transition_immediate, expire_transitioned_object,
|
||||
get_global_expiry_state, get_global_transition_state, init_background_expiry, post_restore_opts,
|
||||
run_stale_multipart_upload_cleanup_once, validate_transition_tier,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -991,6 +991,21 @@ enum ImmediateEnqueueFailure {
|
||||
QueueSendTimedOut { timeout_ms: u64 },
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum TransitionEnqueueOutcome {
|
||||
Queued,
|
||||
AlreadyInFlight,
|
||||
QueueFull,
|
||||
QueueClosed,
|
||||
QueueSendTimedOut,
|
||||
}
|
||||
|
||||
impl TransitionEnqueueOutcome {
|
||||
fn is_handled(self) -> bool {
|
||||
matches!(self, Self::Queued | Self::AlreadyInFlight)
|
||||
}
|
||||
}
|
||||
|
||||
impl TransitionState {
|
||||
#[allow(clippy::new_ret_no_self)]
|
||||
pub fn new() -> Arc<Self> {
|
||||
@@ -1213,12 +1228,17 @@ impl TransitionState {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn queue_transition_task(self: &Arc<Self>, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) -> bool {
|
||||
async fn queue_transition_task_outcome(
|
||||
self: &Arc<Self>,
|
||||
oi: &ObjectInfo,
|
||||
event: &lifecycle::Event,
|
||||
src: &LcEventSrc,
|
||||
) -> TransitionEnqueueOutcome {
|
||||
if is_immediate_transition_source(src) && should_force_immediate_transition_enqueue_timeout() {
|
||||
self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::ForcedTimeout);
|
||||
record_scanner_transition_enqueue_result(src, 1, false);
|
||||
self.record_scanner_transition_state();
|
||||
return false;
|
||||
return TransitionEnqueueOutcome::QueueSendTimedOut;
|
||||
}
|
||||
|
||||
// Deduplicate concurrent enqueues of the same object version. The claim is
|
||||
@@ -1227,7 +1247,7 @@ impl TransitionState {
|
||||
// counters (the first enqueue already counted this object).
|
||||
if !self.reserve_transition(oi) {
|
||||
self.record_scanner_transition_state();
|
||||
return true;
|
||||
return TransitionEnqueueOutcome::AlreadyInFlight;
|
||||
}
|
||||
|
||||
let task = TransitionTask {
|
||||
@@ -1235,15 +1255,14 @@ impl TransitionState {
|
||||
src: src.clone(),
|
||||
event: event.clone(),
|
||||
};
|
||||
let mut queued = false;
|
||||
if is_immediate_transition_source(src) {
|
||||
match self.transition_tx.try_send(Some(task)) {
|
||||
Ok(()) => queued = true,
|
||||
let outcome = match self.transition_tx.try_send(Some(task)) {
|
||||
Ok(()) => TransitionEnqueueOutcome::Queued,
|
||||
Err(async_channel::TrySendError::Full(task)) => {
|
||||
Self::inc_counter(&self.queue_full_tasks);
|
||||
let send_timeout = self.transition_queue_send_timeout;
|
||||
match tokio::time::timeout(send_timeout, self.transition_tx.send(task)).await {
|
||||
Ok(Ok(())) => queued = true,
|
||||
Ok(Ok(())) => TransitionEnqueueOutcome::Queued,
|
||||
Ok(Err(_)) => {
|
||||
self.handle_immediate_enqueue_failure(
|
||||
oi,
|
||||
@@ -1252,6 +1271,7 @@ impl TransitionState {
|
||||
timeout_ms: Some(send_timeout.as_millis() as u64),
|
||||
},
|
||||
);
|
||||
TransitionEnqueueOutcome::QueueClosed
|
||||
}
|
||||
Err(_) => {
|
||||
self.handle_immediate_enqueue_failure(
|
||||
@@ -1261,23 +1281,26 @@ impl TransitionState {
|
||||
timeout_ms: send_timeout.as_millis() as u64,
|
||||
},
|
||||
);
|
||||
TransitionEnqueueOutcome::QueueSendTimedOut
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(async_channel::TrySendError::Closed(_task)) => {
|
||||
self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::QueueClosed { timeout_ms: None });
|
||||
TransitionEnqueueOutcome::QueueClosed
|
||||
}
|
||||
}
|
||||
};
|
||||
let queued = outcome == TransitionEnqueueOutcome::Queued;
|
||||
if !queued {
|
||||
self.release_transition(oi);
|
||||
}
|
||||
record_scanner_transition_enqueue_result(src, 1, queued);
|
||||
self.record_scanner_transition_state();
|
||||
return queued;
|
||||
return outcome;
|
||||
}
|
||||
|
||||
match self.transition_tx.try_send(Some(task)) {
|
||||
Ok(()) => queued = true,
|
||||
let outcome = match self.transition_tx.try_send(Some(task)) {
|
||||
Ok(()) => TransitionEnqueueOutcome::Queued,
|
||||
Err(async_channel::TrySendError::Full(_)) => {
|
||||
Self::inc_counter(&self.queue_full_tasks);
|
||||
debug!(
|
||||
@@ -1286,6 +1309,7 @@ impl TransitionState {
|
||||
source = ?src,
|
||||
"transition queue is full; deferring to scanner/backfill"
|
||||
);
|
||||
TransitionEnqueueOutcome::QueueFull
|
||||
}
|
||||
Err(async_channel::TrySendError::Closed(_)) => {
|
||||
debug!(
|
||||
@@ -1298,14 +1322,20 @@ impl TransitionState {
|
||||
state = "queue_closed",
|
||||
"transition enqueue failed because the queue is closed"
|
||||
);
|
||||
TransitionEnqueueOutcome::QueueClosed
|
||||
}
|
||||
}
|
||||
};
|
||||
let queued = outcome == TransitionEnqueueOutcome::Queued;
|
||||
if !queued {
|
||||
self.release_transition(oi);
|
||||
}
|
||||
record_scanner_transition_enqueue_result(src, 1, queued);
|
||||
self.record_scanner_transition_state();
|
||||
queued
|
||||
outcome
|
||||
}
|
||||
|
||||
pub async fn queue_transition_task(self: &Arc<Self>, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) -> bool {
|
||||
self.queue_transition_task_outcome(oi, event, src).await.is_handled()
|
||||
}
|
||||
|
||||
pub async fn init(api: Arc<ECStore>) {
|
||||
@@ -2402,30 +2432,116 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
|
||||
pub struct ManualTransitionRunOptions {
|
||||
pub prefix: String,
|
||||
pub marker: Option<String>,
|
||||
pub version_marker: Option<String>,
|
||||
pub tier: Option<String>,
|
||||
pub dry_run: bool,
|
||||
pub max_objects: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
|
||||
pub struct ManualTransitionRunReport {
|
||||
pub bucket: String,
|
||||
pub prefix: String,
|
||||
pub tier: Option<String>,
|
||||
pub dry_run: bool,
|
||||
pub lifecycle_config_found: bool,
|
||||
pub scanned: u64,
|
||||
pub eligible: u64,
|
||||
pub enqueued: u64,
|
||||
pub dry_run_eligible: u64,
|
||||
pub skipped_not_transition: u64,
|
||||
pub skipped_tier: u64,
|
||||
pub skipped_delete_marker: u64,
|
||||
pub skipped_directory: u64,
|
||||
pub skipped_replication: u64,
|
||||
pub skipped_already_in_flight: u64,
|
||||
pub skipped_queue_full: u64,
|
||||
pub skipped_queue_closed: u64,
|
||||
pub skipped_queue_timeout: u64,
|
||||
pub truncated_by_limit: bool,
|
||||
#[serde(skip_serializing)]
|
||||
pub next_marker: Option<String>,
|
||||
#[serde(skip_serializing)]
|
||||
pub next_version_idmarker: Option<String>,
|
||||
}
|
||||
|
||||
impl ManualTransitionRunReport {
|
||||
fn new(bucket: &str, options: &ManualTransitionRunOptions) -> Self {
|
||||
Self {
|
||||
bucket: bucket.to_string(),
|
||||
prefix: options.prefix.clone(),
|
||||
tier: options.tier.clone(),
|
||||
dry_run: options.dry_run,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
pub fn has_partial_enqueue(&self) -> bool {
|
||||
self.skipped_queue_full > 0 || self.skipped_queue_closed > 0 || self.skipped_queue_timeout > 0
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn enqueue_transition_for_existing_objects(api: Arc<ECStore>, bucket: &str) -> Result<(), Error> {
|
||||
let _ = enqueue_transition_for_existing_objects_scoped(api, bucket, ManualTransitionRunOptions::default()).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn enqueue_transition_for_existing_objects_scoped(
|
||||
api: Arc<ECStore>,
|
||||
bucket: &str,
|
||||
options: ManualTransitionRunOptions,
|
||||
) -> Result<ManualTransitionRunReport, Error> {
|
||||
const LIST_PAGE_SIZE: i32 = 1000;
|
||||
|
||||
let mut report = ManualTransitionRunReport::new(bucket, &options);
|
||||
let Some(lc) = runtime_sources::bucket_lifecycle_config(bucket).await else {
|
||||
return Ok(());
|
||||
return Ok(report);
|
||||
};
|
||||
let mut marker = None;
|
||||
let mut version_marker = None;
|
||||
report.lifecycle_config_found = true;
|
||||
let mut marker = options.marker.clone();
|
||||
let mut version_marker = options.version_marker.clone();
|
||||
let mut previous_marker = marker.clone();
|
||||
let mut previous_version_marker = version_marker.clone();
|
||||
let src = LcEventSrc::Scanner;
|
||||
|
||||
loop {
|
||||
let page = api
|
||||
.clone()
|
||||
.list_object_versions(bucket, "", marker.clone(), version_marker.clone(), None, 1000)
|
||||
.list_object_versions(bucket, &options.prefix, marker.clone(), version_marker.clone(), None, LIST_PAGE_SIZE)
|
||||
.await?;
|
||||
|
||||
for object in &page.objects {
|
||||
enqueue_transition_with_lifecycle(object, &lc, &src).await;
|
||||
for (index, object) in page.objects.iter().enumerate() {
|
||||
report.scanned = report.scanned.saturating_add(1);
|
||||
enqueue_transition_with_lifecycle_report(object, &lc, &src, &options, &mut report).await;
|
||||
if report.has_partial_enqueue() {
|
||||
report.next_marker.clone_from(&previous_marker);
|
||||
report.next_version_idmarker.clone_from(&previous_version_marker);
|
||||
return Ok(report);
|
||||
}
|
||||
if options.max_objects.is_some_and(|max_objects| report.scanned >= max_objects) {
|
||||
if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) {
|
||||
report.truncated_by_limit = true;
|
||||
report.next_marker = Some(object.name.clone());
|
||||
report.next_version_idmarker = Some(manual_transition_version_marker(object));
|
||||
}
|
||||
return Ok(report);
|
||||
}
|
||||
previous_marker = Some(object.name.clone());
|
||||
previous_version_marker = Some(manual_transition_version_marker(object));
|
||||
}
|
||||
|
||||
if !page.is_truncated {
|
||||
return Ok(());
|
||||
return Ok(report);
|
||||
}
|
||||
|
||||
marker = page.next_marker;
|
||||
version_marker = page.next_version_idmarker;
|
||||
previous_marker = marker.clone();
|
||||
previous_version_marker = version_marker.clone();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2615,6 +2731,76 @@ pub async fn enqueue_expiry_for_existing_objects(api: Arc<ECStore>, bucket: &str
|
||||
}
|
||||
}
|
||||
|
||||
fn manual_transition_has_more_after_limit(page_index: usize, page_len: usize, page_is_truncated: bool) -> bool {
|
||||
page_index.saturating_add(1) < page_len || page_is_truncated
|
||||
}
|
||||
|
||||
fn manual_transition_version_marker(oi: &ObjectInfo) -> String {
|
||||
oi.version_id
|
||||
.map(|version| version.to_string())
|
||||
.unwrap_or_else(|| "null".to_string())
|
||||
}
|
||||
|
||||
async fn enqueue_transition_with_lifecycle_report(
|
||||
oi: &ObjectInfo,
|
||||
lc: &BucketLifecycleConfiguration,
|
||||
src: &LcEventSrc,
|
||||
options: &ManualTransitionRunOptions,
|
||||
report: &mut ManualTransitionRunReport,
|
||||
) -> bool {
|
||||
let event = lc.eval(&oi.to_lifecycle_opts()).await;
|
||||
match event.action {
|
||||
IlmAction::TransitionAction | IlmAction::TransitionVersionAction => {
|
||||
if oi.delete_marker || oi.is_dir {
|
||||
if oi.delete_marker {
|
||||
report.skipped_delete_marker = report.skipped_delete_marker.saturating_add(1);
|
||||
} else {
|
||||
report.skipped_directory = report.skipped_directory.saturating_add(1);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
if lifecycle_action_blocked_by_replication(event.action, oi) {
|
||||
report.skipped_replication = report.skipped_replication.saturating_add(1);
|
||||
return false;
|
||||
}
|
||||
if options
|
||||
.tier
|
||||
.as_deref()
|
||||
.is_some_and(|tier| !event.storage_class.eq_ignore_ascii_case(tier))
|
||||
{
|
||||
report.skipped_tier = report.skipped_tier.saturating_add(1);
|
||||
return false;
|
||||
}
|
||||
report.eligible = report.eligible.saturating_add(1);
|
||||
if options.dry_run {
|
||||
report.dry_run_eligible = report.dry_run_eligible.saturating_add(1);
|
||||
return true;
|
||||
}
|
||||
let outcome = runtime_sources::transition_state_handle()
|
||||
.queue_transition_task_outcome(oi, &event, src)
|
||||
.await;
|
||||
match outcome {
|
||||
TransitionEnqueueOutcome::Queued => report.enqueued = report.enqueued.saturating_add(1),
|
||||
TransitionEnqueueOutcome::AlreadyInFlight => {
|
||||
report.skipped_already_in_flight = report.skipped_already_in_flight.saturating_add(1);
|
||||
}
|
||||
TransitionEnqueueOutcome::QueueFull => {
|
||||
report.skipped_queue_full = report.skipped_queue_full.saturating_add(1);
|
||||
}
|
||||
TransitionEnqueueOutcome::QueueClosed => {
|
||||
report.skipped_queue_closed = report.skipped_queue_closed.saturating_add(1);
|
||||
}
|
||||
TransitionEnqueueOutcome::QueueSendTimedOut => {
|
||||
report.skipped_queue_timeout = report.skipped_queue_timeout.saturating_add(1);
|
||||
}
|
||||
}
|
||||
return outcome.is_handled();
|
||||
}
|
||||
_ => report.skipped_not_transition = report.skipped_not_transition.saturating_add(1),
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
async fn enqueue_transition_with_lifecycle(oi: &ObjectInfo, lc: &BucketLifecycleConfiguration, src: &LcEventSrc) -> bool {
|
||||
let event = lc.eval(&oi.to_lifecycle_opts()).await;
|
||||
match event.action {
|
||||
@@ -2625,13 +2811,12 @@ async fn enqueue_transition_with_lifecycle(oi: &ObjectInfo, lc: &BucketLifecycle
|
||||
if lifecycle_action_blocked_by_replication(event.action, oi) {
|
||||
return false;
|
||||
}
|
||||
return runtime_sources::transition_state_handle()
|
||||
runtime_sources::transition_state_handle()
|
||||
.queue_transition_task(oi, &event, src)
|
||||
.await;
|
||||
.await
|
||||
}
|
||||
_ => (),
|
||||
_ => false,
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
/// Build the delete options for a lifecycle expiry event on a transitioned
|
||||
@@ -3459,13 +3644,15 @@ pub async fn apply_lifecycle_action(event: &lifecycle::Event, src: &LcEventSrc,
|
||||
mod tests {
|
||||
use super::{
|
||||
DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS, DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX,
|
||||
DEFAULT_TRANSITION_WORKERS_CAP, ExpiryState, FreeVersionTask, StaleMultipartUploadCandidate,
|
||||
TIER_FREE_VERSION_RECOVERY_BASE_INTERVAL, TIER_FREE_VERSION_RECOVERY_MAX_IDLE_INTERVAL, TierFreeVersionRecoverySchedule,
|
||||
TransitionState, TransitionedObject, VersionReplicationScan, cleanup_empty_multipart_sha_dirs_on_local_disks,
|
||||
cleanup_stale_multipart_uploads_once_at, enqueue_recovered_free_version_with_state, enqueue_transition_with_lifecycle,
|
||||
DEFAULT_TRANSITION_WORKERS_CAP, ExpiryState, FreeVersionTask, ManualTransitionRunOptions, ManualTransitionRunReport,
|
||||
StaleMultipartUploadCandidate, TIER_FREE_VERSION_RECOVERY_BASE_INTERVAL, TIER_FREE_VERSION_RECOVERY_MAX_IDLE_INTERVAL,
|
||||
TierFreeVersionRecoverySchedule, TransitionEnqueueOutcome, TransitionState, TransitionedObject, VersionReplicationScan,
|
||||
cleanup_empty_multipart_sha_dirs_on_local_disks, cleanup_stale_multipart_uploads_once_at,
|
||||
enqueue_recovered_free_version_with_state, enqueue_transition_with_lifecycle, enqueue_transition_with_lifecycle_report,
|
||||
eval_action_from_lifecycle, jitter_tier_free_version_recovery_delay, lifecycle_action_blocked_by_replication,
|
||||
lifecycle_delete_all_versions_replication_scan, lifecycle_deleted_object, lifecycle_replication_blocks_action,
|
||||
lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets,
|
||||
manual_transition_has_more_after_limit, manual_transition_version_marker,
|
||||
mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, replication_state_for_delete,
|
||||
resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count,
|
||||
resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location,
|
||||
@@ -5594,6 +5781,63 @@ mod tests {
|
||||
assert_eq!(state.transition_rx.len(), 1, "only one task must actually be queued");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn queue_transition_task_outcome_reports_duplicate_separately() {
|
||||
let state = TransitionState::new_with_capacity(4);
|
||||
let object = ObjectInfo {
|
||||
bucket: "bucket".to_string(),
|
||||
name: "object".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
let event = crate::bucket::lifecycle::lifecycle::Event {
|
||||
action: IlmAction::TransitionAction,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let first = state
|
||||
.queue_transition_task_outcome(&object, &event, &LcEventSrc::Scanner)
|
||||
.await;
|
||||
let second = state
|
||||
.queue_transition_task_outcome(&object, &event, &LcEventSrc::Scanner)
|
||||
.await;
|
||||
|
||||
assert_eq!(first, TransitionEnqueueOutcome::Queued);
|
||||
assert_eq!(second, TransitionEnqueueOutcome::AlreadyInFlight);
|
||||
assert_eq!(state.transition_rx.len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn queue_transition_task_outcome_reports_queue_full_separately() {
|
||||
let state = TransitionState::new_with_capacity(1);
|
||||
let first_object = ObjectInfo {
|
||||
bucket: "bucket".to_string(),
|
||||
name: "first".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
let second_object = ObjectInfo {
|
||||
bucket: "bucket".to_string(),
|
||||
name: "second".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
let event = crate::bucket::lifecycle::lifecycle::Event {
|
||||
action: IlmAction::TransitionAction,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let first = state
|
||||
.queue_transition_task_outcome(&first_object, &event, &LcEventSrc::Scanner)
|
||||
.await;
|
||||
let second = state
|
||||
.queue_transition_task_outcome(&second_object, &event, &LcEventSrc::Scanner)
|
||||
.await;
|
||||
|
||||
assert_eq!(first, TransitionEnqueueOutcome::Queued);
|
||||
assert_eq!(second, TransitionEnqueueOutcome::QueueFull);
|
||||
assert_eq!(state.transition_rx.len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn queue_transition_task_dedupes_immediate_and_scanner_sources_for_same_version() {
|
||||
@@ -6031,6 +6275,85 @@ mod tests {
|
||||
assert!(!queued);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn manual_transition_dry_run_counts_due_object_without_enqueue() {
|
||||
let lc = latest_transition_lifecycle();
|
||||
let object = current_object(ReplicationStatusType::Completed);
|
||||
let options = ManualTransitionRunOptions {
|
||||
dry_run: true,
|
||||
..Default::default()
|
||||
};
|
||||
let mut report = ManualTransitionRunReport::new(&object.bucket, &options);
|
||||
|
||||
let handled = enqueue_transition_with_lifecycle_report(&object, &lc, &LcEventSrc::Scanner, &options, &mut report).await;
|
||||
|
||||
assert!(handled);
|
||||
assert_eq!(report.eligible, 1);
|
||||
assert_eq!(report.dry_run_eligible, 1);
|
||||
assert_eq!(report.enqueued, 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn manual_transition_respects_not_yet_due_lifecycle_rule() {
|
||||
let lc = latest_transition_lifecycle();
|
||||
let object = ObjectInfo {
|
||||
mod_time: Some(OffsetDateTime::now_utc()),
|
||||
..current_object(ReplicationStatusType::Completed)
|
||||
};
|
||||
let options = ManualTransitionRunOptions {
|
||||
dry_run: true,
|
||||
..Default::default()
|
||||
};
|
||||
let mut report = ManualTransitionRunReport::new(&object.bucket, &options);
|
||||
|
||||
let handled = enqueue_transition_with_lifecycle_report(&object, &lc, &LcEventSrc::Scanner, &options, &mut report).await;
|
||||
|
||||
assert!(!handled);
|
||||
assert_eq!(report.eligible, 0);
|
||||
assert_eq!(report.skipped_not_transition, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn manual_transition_tier_filter_skips_non_matching_transition() {
|
||||
let lc = latest_transition_lifecycle();
|
||||
let object = current_object(ReplicationStatusType::Completed);
|
||||
let options = ManualTransitionRunOptions {
|
||||
tier: Some("COLD".to_string()),
|
||||
dry_run: true,
|
||||
..Default::default()
|
||||
};
|
||||
let mut report = ManualTransitionRunReport::new(&object.bucket, &options);
|
||||
|
||||
let handled = enqueue_transition_with_lifecycle_report(&object, &lc, &LcEventSrc::Scanner, &options, &mut report).await;
|
||||
|
||||
assert!(!handled);
|
||||
assert_eq!(report.eligible, 0);
|
||||
assert_eq!(report.skipped_tier, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_version_marker_preserves_null_version_cursor() {
|
||||
let null_version = ObjectInfo {
|
||||
version_id: None,
|
||||
..Default::default()
|
||||
};
|
||||
let version_id = Uuid::new_v4();
|
||||
let versioned = ObjectInfo {
|
||||
version_id: Some(version_id),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert_eq!(manual_transition_version_marker(&null_version), "null");
|
||||
assert_eq!(manual_transition_version_marker(&versioned), version_id.to_string());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_limit_boundary_only_truncates_when_more_objects_exist() {
|
||||
assert!(!manual_transition_has_more_after_limit(9, 10, false));
|
||||
assert!(manual_transition_has_more_after_limit(9, 11, false));
|
||||
assert!(manual_transition_has_more_after_limit(9, 10, true));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() {
|
||||
let lc = expired_delete_marker_lifecycle();
|
||||
|
||||
@@ -0,0 +1,256 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crate::admin::auth::validate_admin_request;
|
||||
use crate::admin::router::{AdminOperation, Operation, S3Router};
|
||||
use crate::admin::runtime_sources::object_store_from_extensions;
|
||||
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
|
||||
use crate::admin::storage_api::lifecycle::{
|
||||
ManualTransitionRunOptions, ManualTransitionRunReport, enqueue_transition_for_existing_objects_scoped,
|
||||
};
|
||||
use crate::auth::{check_key_valid, get_session_token};
|
||||
use crate::server::{ADMIN_PREFIX, RemoteAddr};
|
||||
use http::{HeaderMap, HeaderValue};
|
||||
use hyper::{Method, StatusCode};
|
||||
use matchit::Params;
|
||||
use rustfs_policy::policy::action::{Action, AdminAction};
|
||||
use s3s::header::CONTENT_TYPE;
|
||||
use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
const JSON_CONTENT_TYPE: &str = "application/json";
|
||||
const DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS: u64 = 10_000;
|
||||
const MAX_MANUAL_TRANSITION_OBJECTS: u64 = 100_000;
|
||||
|
||||
#[derive(Debug, Deserialize, Default)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct ManualTransitionRunQuery {
|
||||
bucket: Option<String>,
|
||||
prefix: Option<String>,
|
||||
marker: Option<String>,
|
||||
#[serde(rename = "versionMarker")]
|
||||
version_marker: Option<String>,
|
||||
tier: Option<String>,
|
||||
#[serde(rename = "dryRun")]
|
||||
dry_run: Option<bool>,
|
||||
#[serde(rename = "maxObjects")]
|
||||
max_objects: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
struct ManualTransitionRunResponse {
|
||||
state: &'static str,
|
||||
mode: &'static str,
|
||||
job_id: Option<String>,
|
||||
status_endpoint: Option<String>,
|
||||
report: ManualTransitionRunReport,
|
||||
}
|
||||
|
||||
pub fn register_ilm_transition_route(r: &mut S3Router<AdminOperation>) -> std::io::Result<()> {
|
||||
r.insert(
|
||||
Method::POST,
|
||||
format!("{ADMIN_PREFIX}/v3/ilm/transition/run").as_str(),
|
||||
AdminOperation(&ManualTransitionRunHandler {}),
|
||||
)?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, ManualTransitionRunOptions)> {
|
||||
let query: ManualTransitionRunQuery = match query {
|
||||
Some(query) => serde_urlencoded::from_bytes(query.as_bytes())
|
||||
.map_err(|_| s3_error!(InvalidArgument, "invalid manual transition query"))?,
|
||||
None => ManualTransitionRunQuery::default(),
|
||||
};
|
||||
|
||||
let bucket = query
|
||||
.bucket
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|bucket| !bucket.is_empty())
|
||||
.ok_or_else(|| s3_error!(InvalidRequest, "bucket is required"))?;
|
||||
if is_reserved_or_invalid_bucket(bucket, false) {
|
||||
return Err(s3_error!(InvalidBucketName, "invalid bucket name"));
|
||||
}
|
||||
|
||||
let max_objects = query.max_objects.unwrap_or(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS);
|
||||
if max_objects == 0 || max_objects > MAX_MANUAL_TRANSITION_OBJECTS {
|
||||
return Err(s3_error!(InvalidArgument, "maxObjects is outside the allowed range"));
|
||||
}
|
||||
|
||||
Ok((
|
||||
bucket.to_string(),
|
||||
ManualTransitionRunOptions {
|
||||
prefix: query.prefix.unwrap_or_default(),
|
||||
marker: query.marker.filter(|marker| !marker.is_empty()),
|
||||
version_marker: query.version_marker.filter(|version_marker| !version_marker.is_empty()),
|
||||
tier: query.tier.map(|tier| tier.trim().to_string()).filter(|tier| !tier.is_empty()),
|
||||
dry_run: query.dry_run.unwrap_or(false),
|
||||
max_objects: Some(max_objects),
|
||||
},
|
||||
))
|
||||
}
|
||||
|
||||
async fn authorize_manual_transition_request(req: &S3Request<Body>) -> S3Result<()> {
|
||||
let Some(input_cred) = req.credentials.as_ref() else {
|
||||
return Err(s3_error!(InvalidRequest, "authentication required"));
|
||||
};
|
||||
|
||||
let (cred, owner) =
|
||||
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?;
|
||||
let remote_addr = req
|
||||
.extensions
|
||||
.get::<Option<RemoteAddr>>()
|
||||
.and_then(|opt| opt.map(|addr| addr.0));
|
||||
|
||||
validate_admin_request(
|
||||
&req.headers,
|
||||
&cred,
|
||||
owner,
|
||||
false,
|
||||
vec![Action::AdminAction(AdminAction::SetTierAction)],
|
||||
remote_addr,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
fn response_state(report: &ManualTransitionRunReport) -> &'static str {
|
||||
if report.truncated_by_limit || report.has_partial_enqueue() {
|
||||
"partial"
|
||||
} else {
|
||||
"completed"
|
||||
}
|
||||
}
|
||||
|
||||
fn json_response(response: &ManualTransitionRunResponse) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
let body = serde_json::to_vec(response).map_err(|err| {
|
||||
S3Error::with_message(S3ErrorCode::InternalError, format!("failed to encode manual transition response: {err}"))
|
||||
})?;
|
||||
let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE)
|
||||
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("invalid content type: {err}")))?;
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert(CONTENT_TYPE, content_type);
|
||||
Ok(S3Response::with_headers((StatusCode::OK, Body::from(body)), headers))
|
||||
}
|
||||
|
||||
pub struct ManualTransitionRunHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for ManualTransitionRunHandler {
|
||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
authorize_manual_transition_request(&req).await?;
|
||||
let (bucket, options) = parse_manual_transition_query(req.uri.query())?;
|
||||
let Some(store) = object_store_from_extensions(&req.extensions) else {
|
||||
return Err(s3_error!(InternalError, "object store is not initialized"));
|
||||
};
|
||||
|
||||
let report = enqueue_transition_for_existing_objects_scoped(store, &bucket, options)
|
||||
.await
|
||||
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("manual transition run failed: {err}")))?;
|
||||
let response = ManualTransitionRunResponse {
|
||||
state: response_state(&report),
|
||||
mode: "enqueue_only",
|
||||
job_id: None,
|
||||
status_endpoint: None,
|
||||
report,
|
||||
};
|
||||
|
||||
json_response(&response)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn manual_transition_query_defaults_to_bounded_run() {
|
||||
let (bucket, options) =
|
||||
parse_manual_transition_query(Some("bucket=data&prefix=logs/&marker=logs/a&versionMarker=v1&tier=warm"))
|
||||
.expect("valid query should parse");
|
||||
|
||||
assert_eq!(bucket, "data");
|
||||
assert_eq!(options.prefix, "logs/");
|
||||
assert_eq!(options.marker.as_deref(), Some("logs/a"));
|
||||
assert_eq!(options.version_marker.as_deref(), Some("v1"));
|
||||
assert_eq!(options.tier.as_deref(), Some("warm"));
|
||||
assert!(!options.dry_run);
|
||||
assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_query_rejects_server_info_style_unscoped_request() {
|
||||
let err = parse_manual_transition_query(Some("dryRun=true")).expect_err("bucket must be required");
|
||||
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_query_rejects_unbounded_budget() {
|
||||
let err = parse_manual_transition_query(Some("bucket=data&maxObjects=0")).expect_err("zero budget must fail");
|
||||
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_response_reports_partial_for_queue_pressure() {
|
||||
let report = ManualTransitionRunReport {
|
||||
skipped_queue_full: 1,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert_eq!(response_state(&report), "partial");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_response_omits_raw_resume_markers() {
|
||||
let report = ManualTransitionRunReport {
|
||||
truncated_by_limit: true,
|
||||
next_marker: Some("private/object".to_string()),
|
||||
next_version_idmarker: Some("null".to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
let response = ManualTransitionRunResponse {
|
||||
state: response_state(&report),
|
||||
mode: "enqueue_only",
|
||||
job_id: None,
|
||||
status_endpoint: None,
|
||||
report,
|
||||
};
|
||||
|
||||
let value = serde_json::to_value(response).expect("response should serialize");
|
||||
assert!(value.pointer("/report/next_marker").is_none());
|
||||
assert!(value.pointer("/report/next_version_idmarker").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_handler_requires_set_tier_action() {
|
||||
let src = include_str!("ilm_transition.rs");
|
||||
let auth_block = extract_block_between_markers(src, "async fn authorize_manual_transition_request", "fn response_state");
|
||||
|
||||
assert!(auth_block.contains("AdminAction::SetTierAction"));
|
||||
assert!(!auth_block.contains("AdminAction::ServerInfoAdminAction"));
|
||||
}
|
||||
|
||||
fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str {
|
||||
let start = src
|
||||
.find(start_marker)
|
||||
.unwrap_or_else(|| panic!("expected start marker `{start_marker}`"));
|
||||
let after_start = &src[start..];
|
||||
let end = after_start
|
||||
.find(end_marker)
|
||||
.unwrap_or_else(|| panic!("expected end marker `{end_marker}` after `{start_marker}`"));
|
||||
&after_start[..end]
|
||||
}
|
||||
}
|
||||
@@ -28,6 +28,7 @@ pub mod heal;
|
||||
pub mod health;
|
||||
pub(crate) mod iam_error;
|
||||
pub mod idp_compat;
|
||||
pub mod ilm_transition;
|
||||
pub mod is_admin;
|
||||
pub mod kms;
|
||||
pub mod kms_dynamic;
|
||||
@@ -118,6 +119,7 @@ mod tests {
|
||||
let _list_remote_target_handler = replication::ListRemoteTargetHandler {};
|
||||
let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {};
|
||||
let _scanner_status_handler = scanner::ScannerStatusHandler {};
|
||||
let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {};
|
||||
let _site_replication_add_handler = site_replication::SiteReplicationAddHandler {};
|
||||
let _site_replication_info_handler = site_replication::SiteReplicationInfoHandler {};
|
||||
let _site_replication_status_handler = site_replication::SiteReplicationStatusHandler {};
|
||||
|
||||
@@ -34,7 +34,7 @@ mod route_registration_test;
|
||||
|
||||
use handlers::{
|
||||
audit, batch_job, bucket_meta, cluster_snapshot, config_admin, diagnostics, durability as durability_handler, extensions,
|
||||
heal, health, idp_compat, kms, module_switch, object_data_cache, object_zip_download, oidc, plugins_catalog,
|
||||
heal, health, idp_compat, ilm_transition, kms, module_switch, object_data_cache, object_zip_download, oidc, plugins_catalog,
|
||||
plugins_instances, pools, profile_admin, quota as quota_handler, rebalance, replication as replication_handler, scanner,
|
||||
site_replication, sts, system, table_catalog, tier, tls_debug, user,
|
||||
};
|
||||
@@ -76,6 +76,7 @@ fn register_admin_routes(r: &mut S3Router<AdminOperation>) -> std::io::Result<()
|
||||
bucket_meta::register_bucket_meta_route(r)?;
|
||||
config_admin::register_config_route(r)?;
|
||||
scanner::register_scanner_route(r)?;
|
||||
ilm_transition::register_ilm_transition_route(r)?;
|
||||
object_data_cache::register_object_data_cache_route(r)?;
|
||||
audit::register_audit_target_route(r)?;
|
||||
module_switch::register_module_switch_route(r)?;
|
||||
|
||||
@@ -412,6 +412,7 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
|
||||
admin(HttpMethod::Get, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High),
|
||||
admin(HttpMethod::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High),
|
||||
admin(HttpMethod::Get, "/rustfs/admin/v3/scanner/status", SERVER_INFO, RouteRiskLevel::Sensitive),
|
||||
admin(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER, RouteRiskLevel::High),
|
||||
admin(
|
||||
HttpMethod::Get,
|
||||
"/rustfs/admin/v3/audit/target/list",
|
||||
@@ -1791,6 +1792,12 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn route_policy_requires_set_tier_for_manual_transition_run() {
|
||||
assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER);
|
||||
assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SERVER_INFO);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn route_policy_keeps_contextual_auth_deferred() {
|
||||
assert_deferred(
|
||||
|
||||
@@ -194,6 +194,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
|
||||
admin_route(Method::PUT, "/v3/tier"),
|
||||
admin_route_sample(Method::POST, "/v3/tier/{tiername}", "/v3/tier/HOT"),
|
||||
admin_route(Method::POST, "/v3/tier/clear"),
|
||||
admin_route(Method::POST, "/v3/ilm/transition/run"),
|
||||
admin_route(Method::PUT, "/v3/set-bucket-quota"),
|
||||
admin_route(Method::GET, "/v3/get-bucket-quota"),
|
||||
admin_route_sample(Method::PUT, "/v3/quota/{bucket}", "/v3/quota/test-bucket"),
|
||||
@@ -830,6 +831,7 @@ fn test_register_routes_cover_representative_admin_paths() {
|
||||
assert_route(&router, Method::GET, &admin_path("/v3/config"));
|
||||
assert_route(&router, Method::PUT, &admin_path("/v3/config"));
|
||||
assert_route(&router, Method::GET, &admin_path("/v3/scanner/status"));
|
||||
assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/run"));
|
||||
|
||||
assert_route(&router, Method::GET, &table_catalog_path("/config"));
|
||||
assert_route(&router, Method::PUT, &table_catalog_path("/buckets/analytics"));
|
||||
|
||||
@@ -193,6 +193,21 @@ pub(crate) mod bucket_target_sys {
|
||||
}
|
||||
|
||||
pub(crate) mod lifecycle {
|
||||
pub(crate) type ManualTransitionRunOptions =
|
||||
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunOptions;
|
||||
pub(crate) type ManualTransitionRunReport = super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::ManualTransitionRunReport;
|
||||
|
||||
pub(crate) async fn enqueue_transition_for_existing_objects_scoped(
|
||||
api: std::sync::Arc<super::ECStore>,
|
||||
bucket: &str,
|
||||
options: ManualTransitionRunOptions,
|
||||
) -> super::Result<ManualTransitionRunReport> {
|
||||
super::ecstore_bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects_scoped(
|
||||
api, bucket, options,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) mod tier_last_day_stats {
|
||||
#[cfg(test)]
|
||||
pub(crate) type LastDayTierStats = super::super::ecstore_bucket::lifecycle::tier_last_day_stats::LastDayTierStats;
|
||||
|
||||
Reference in New Issue
Block a user