mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-16 09:58:21 +00:00
fix(ecstore): reconcile object cleanup receipts (#6077)
* fix(s3): keep multipart completion publication owned Co-Authored-By: heihutu <heihutu@gmail.com> * fix(s3): keep put publication owned Co-Authored-By: heihutu <heihutu@gmail.com> * chore(app): route multipart context through facade Co-Authored-By: heihutu <heihutu@gmail.com> * fix(ecstore): gate object transaction fencing Co-Authored-By: heihutu <heihutu@gmail.com> * fix(ecstore): fence object transaction epochs Co-Authored-By: heihutu <heihutu@gmail.com> * fix(ecstore): reconcile old data cleanup receipts Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -46,6 +46,7 @@ use super::storage_api::multipart_usecase::options::{
|
||||
get_content_sha256_with_query, get_opts, namespace_reserved_user_metadata, parse_copy_source_range,
|
||||
put_opts_with_replication_authorization, validate_archive_content_encoding,
|
||||
};
|
||||
use super::storage_api::multipart_usecase::request_context::spawn_traced_join;
|
||||
use super::storage_api::multipart_usecase::s3_api::multipart::{
|
||||
ListMultipartUploadsParams, build_list_multipart_uploads_output, build_list_parts_output,
|
||||
parse_list_multipart_uploads_params, parse_list_parts_params, parse_upload_part_number,
|
||||
@@ -588,56 +589,94 @@ impl DefaultMultipartUsecase {
|
||||
None => None,
|
||||
};
|
||||
|
||||
let obj_info = store
|
||||
.clone()
|
||||
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, &opts)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
let _ = invalidate_object_data_cache_after_complete_multipart_success(&cache_adapter, &bucket, &key).await;
|
||||
record_capacity_write(Some(capacity_scope_token)).await;
|
||||
|
||||
if let Some(metadata_sys) = quota_metadata_sys.as_ref() {
|
||||
if opts.replication_request {
|
||||
let quota_checker = QuotaChecker::new(metadata_sys.clone());
|
||||
match quota_checker
|
||||
.check_quota(&bucket, QuotaOperation::PutObject, obj_info.size.max(0) as u64)
|
||||
let complete_commit = spawn_traced_join({
|
||||
let store = Arc::clone(&store);
|
||||
let bucket = bucket.clone();
|
||||
let key = key.clone();
|
||||
let upload_id = upload_id.clone();
|
||||
let opts = opts.clone();
|
||||
let quota_metadata_sys = quota_metadata_sys.clone();
|
||||
async move {
|
||||
let obj_info = store
|
||||
.clone()
|
||||
.complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, &opts)
|
||||
.await
|
||||
{
|
||||
Ok(check_result) if !check_result.allowed => {
|
||||
let _ = store.delete_object(&bucket, &key, ObjectOptions::default()).await;
|
||||
let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await;
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::InvalidRequest,
|
||||
format!(
|
||||
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
|
||||
check_result.current_usage.unwrap_or(0),
|
||||
check_result.quota_limit.unwrap_or(0)
|
||||
),
|
||||
));
|
||||
.map_err(ApiError::from)?;
|
||||
let _ = invalidate_object_data_cache_after_complete_multipart_success(&cache_adapter, &bucket, &key).await;
|
||||
record_capacity_write(Some(capacity_scope_token)).await;
|
||||
|
||||
if let Some(metadata_sys) = quota_metadata_sys.as_ref() {
|
||||
if opts.replication_request {
|
||||
let quota_checker = QuotaChecker::new(metadata_sys.clone());
|
||||
match quota_checker
|
||||
.check_quota(&bucket, QuotaOperation::PutObject, obj_info.size.max(0) as u64)
|
||||
.await
|
||||
{
|
||||
Ok(check_result) if !check_result.allowed => {
|
||||
let _ = store.delete_object(&bucket, &key, ObjectOptions::default()).await;
|
||||
let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await;
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::InvalidRequest,
|
||||
format!(
|
||||
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
|
||||
check_result.current_usage.unwrap_or(0),
|
||||
check_result.quota_limit.unwrap_or(0)
|
||||
),
|
||||
));
|
||||
}
|
||||
Err(err) => {
|
||||
warn!("Quota check failed for bucket {} after multipart completion: {}", bucket, err);
|
||||
}
|
||||
Ok(_) => {}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
warn!("Quota check failed for bucket {} after multipart completion: {}", bucket, err);
|
||||
|
||||
let committed_size = if opts.replication_request {
|
||||
obj_info.size.max(0) as u64
|
||||
} else {
|
||||
quota_accounting_object_size(&obj_info, opts.quota_admission.is_some())?
|
||||
};
|
||||
if versioned {
|
||||
record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await;
|
||||
} else {
|
||||
record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await;
|
||||
}
|
||||
Ok(_) => {}
|
||||
}
|
||||
|
||||
enqueue_transition_immediate(&obj_info, LcEventSrc::S3CompleteMultipartUpload).await;
|
||||
|
||||
let mt2 = obj_info.user_defined.clone();
|
||||
let dsc = must_replicate_object(
|
||||
&bucket,
|
||||
&key,
|
||||
&mt2,
|
||||
"".to_string(),
|
||||
opts.delete_marker_replication_status(),
|
||||
opts.clone(),
|
||||
)
|
||||
.await;
|
||||
|
||||
if dsc.replicate_any() {
|
||||
warn!("need multipart replication");
|
||||
schedule_object_replication(obj_info.clone(), store, dsc).await;
|
||||
}
|
||||
|
||||
rustfs_scanner::record_dirty_usage_bucket(&bucket);
|
||||
Ok::<_, S3Error>(obj_info)
|
||||
}
|
||||
});
|
||||
let obj_info = complete_commit.await.map_err(|err| {
|
||||
S3Error::with_message(
|
||||
S3ErrorCode::InternalError,
|
||||
format!("complete multipart upload commit owner task failed: {err}"),
|
||||
)
|
||||
})??;
|
||||
|
||||
let committed_size = if opts.replication_request {
|
||||
obj_info.size.max(0) as u64
|
||||
} else {
|
||||
quota_accounting_object_size(&obj_info, opts.quota_admission.is_some())?
|
||||
};
|
||||
if versioned {
|
||||
record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await;
|
||||
} else {
|
||||
record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await;
|
||||
}
|
||||
}
|
||||
|
||||
enqueue_transition_immediate(&obj_info, LcEventSrc::S3CompleteMultipartUpload).await;
|
||||
|
||||
let raw_mpu_version = obj_info.version_id.map(|v| v.to_string());
|
||||
let mpu_version = if versioned { raw_mpu_version.clone() } else { None };
|
||||
let mpu_version = if versioned {
|
||||
obj_info.version_id.map(|v| v.to_string())
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let mpu_version_for_event = mpu_version.clone();
|
||||
// checksum: stored (decrypted) values take precedence over the request input;
|
||||
// additional algorithms (XXHash3/64/128, SHA-512, MD5), which have no typed
|
||||
@@ -660,28 +699,18 @@ impl DefaultMultipartUsecase {
|
||||
bucket: Some(bucket.clone()),
|
||||
key: Some(key.clone()),
|
||||
e_tag: obj_info.etag.clone().map(|etag| to_s3s_etag(&etag)),
|
||||
location: Some(location.clone()),
|
||||
location: Some(location),
|
||||
server_side_encryption: server_side_encryption.clone(),
|
||||
ssekms_key_id: ssekms_key_id.clone(),
|
||||
checksum_crc32: checksum_crc32.clone(),
|
||||
checksum_crc32c: checksum_crc32c.clone(),
|
||||
checksum_sha1: checksum_sha1.clone(),
|
||||
checksum_sha256: checksum_sha256.clone(),
|
||||
checksum_crc64nvme: checksum_crc64nvme.clone(),
|
||||
checksum_type: checksum_type.clone(),
|
||||
checksum_crc32,
|
||||
checksum_crc32c,
|
||||
checksum_sha1,
|
||||
checksum_sha256,
|
||||
checksum_crc64nvme,
|
||||
checksum_type,
|
||||
version_id: mpu_version,
|
||||
..Default::default()
|
||||
};
|
||||
let mt2 = obj_info.user_defined.clone();
|
||||
let dsc =
|
||||
must_replicate_object(&bucket, &key, &mt2, "".to_string(), opts.delete_marker_replication_status(), opts.clone())
|
||||
.await;
|
||||
|
||||
if dsc.replicate_any() {
|
||||
warn!("need multipart replication");
|
||||
schedule_object_replication(obj_info.clone(), store, dsc).await;
|
||||
}
|
||||
|
||||
// Set object info for event notification
|
||||
helper = helper.object(obj_info);
|
||||
if let Some(version_id) = &mpu_version_for_event {
|
||||
@@ -712,7 +741,6 @@ impl DefaultMultipartUsecase {
|
||||
}
|
||||
let result = Ok(response);
|
||||
let _ = helper.complete(&result);
|
||||
rustfs_scanner::record_dirty_usage_bucket(&bucket);
|
||||
result
|
||||
}
|
||||
|
||||
|
||||
+208
-104
@@ -2989,6 +2989,11 @@ struct PutObjectChecksums {
|
||||
crc64nvme: Option<String>,
|
||||
}
|
||||
|
||||
struct PutObjectCommitResult {
|
||||
obj_info: ObjectInfo,
|
||||
put_versioned: bool,
|
||||
}
|
||||
|
||||
fn normalize_delete_objects_version_id(
|
||||
version_id: Option<String>,
|
||||
) -> std::result::Result<(Option<String>, Option<Uuid>), String> {
|
||||
@@ -5932,7 +5937,7 @@ impl DefaultObjectUsecase {
|
||||
reader = write_plan.apply(reader, actual_size).map_err(ApiError::from)?;
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_encryption_prepare", encryption_stage_start);
|
||||
|
||||
let mut reader = PutObjReader::new(reader);
|
||||
let reader = PutObjReader::new(reader);
|
||||
|
||||
let mt2 = metadata.clone();
|
||||
opts.user_defined.extend(metadata);
|
||||
@@ -6005,97 +6010,145 @@ impl DefaultObjectUsecase {
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let object_traffic_progress = object_traffic_health
|
||||
.as_deref()
|
||||
.and_then(ObjectTrafficHealth::track_write_storage);
|
||||
let store_put_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
let (obj_info, backfilled_old_current_size) = match store
|
||||
.put_object_with_old_current_size(&bucket, &key, &mut reader, &opts)
|
||||
.await
|
||||
.map_err(ApiError::from)
|
||||
{
|
||||
Ok(obj_info) => {
|
||||
store_put_watchdog.cancel();
|
||||
debug!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
event = EVENT_PUT_OBJECT_STORE_RETURNED,
|
||||
component = LOG_COMPONENT_APP,
|
||||
subsystem = LOG_SUBSYSTEM_OBJECT,
|
||||
request_id = %request_id,
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_path = put_path,
|
||||
object_size = actual_size,
|
||||
duration_ms = start_time.elapsed().as_millis() as u64,
|
||||
result = "success",
|
||||
"PutObject store write returned"
|
||||
);
|
||||
obj_info
|
||||
let put_commit = spawn_traced_join({
|
||||
let store = Arc::clone(&store);
|
||||
let bucket = bucket.clone();
|
||||
let key = key.clone();
|
||||
let opts = opts.clone();
|
||||
let cache_adapter = cache_adapter.clone();
|
||||
let request_id = request_id.clone();
|
||||
let put_path = put_path.to_string();
|
||||
async move {
|
||||
let object_traffic_progress = object_traffic_health
|
||||
.as_deref()
|
||||
.and_then(ObjectTrafficHealth::track_write_storage);
|
||||
let mut reader = reader;
|
||||
let store_put_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
let (obj_info, backfilled_old_current_size) = match store
|
||||
.put_object_with_old_current_size(&bucket, &key, &mut reader, &opts)
|
||||
.await
|
||||
.map_err(ApiError::from)
|
||||
{
|
||||
Ok(obj_info) => {
|
||||
store_put_watchdog.cancel();
|
||||
debug!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
event = EVENT_PUT_OBJECT_STORE_RETURNED,
|
||||
component = LOG_COMPONENT_APP,
|
||||
subsystem = LOG_SUBSYSTEM_OBJECT,
|
||||
request_id = %request_id,
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_path = %put_path,
|
||||
object_size = actual_size,
|
||||
duration_ms = start_time.elapsed().as_millis() as u64,
|
||||
result = "success",
|
||||
"PutObject store write returned"
|
||||
);
|
||||
obj_info
|
||||
}
|
||||
Err(err) => {
|
||||
store_put_watchdog.cancel();
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start);
|
||||
warn!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
event = EVENT_PUT_OBJECT_STORE_RETURNED,
|
||||
component = LOG_COMPONENT_APP,
|
||||
subsystem = LOG_SUBSYSTEM_OBJECT,
|
||||
request_id = %request_id,
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_path = %put_path,
|
||||
object_size = actual_size,
|
||||
duration_ms = start_time.elapsed().as_millis() as u64,
|
||||
result = "error",
|
||||
error = %err,
|
||||
"PutObject store write returned"
|
||||
);
|
||||
return Err(err.into());
|
||||
}
|
||||
};
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start);
|
||||
drop(object_traffic_progress);
|
||||
#[cfg(test)]
|
||||
wait_for_put_post_store_test_hook(&bucket).await;
|
||||
|
||||
let post_store_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
maybe_enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await;
|
||||
let _ = invalidate_object_data_cache_after_put_success(&cache_adapter, &bucket, &key).await;
|
||||
|
||||
let put_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
|
||||
// Fast in-memory update for immediate quota and admin usage consistency.
|
||||
// The previous current size comes from the prelookup when it ran,
|
||||
// otherwise from the rename_data backfill (rustfs/backlog#1009); the
|
||||
// backfill reproduces the lookup's observation bit for bit (latest
|
||||
// version's ObjectInfo.size — 0 for a delete-marker latest — or
|
||||
// not-found → None).
|
||||
match prelookup_previous_current_size.or_else(|| previous_current_size_from_backfill(backfilled_old_current_size))
|
||||
{
|
||||
Some(previous_current_size) => {
|
||||
if put_versioned {
|
||||
record_bucket_object_version_write_memory(
|
||||
&bucket,
|
||||
previous_current_size,
|
||||
obj_info.size.max(0) as u64,
|
||||
)
|
||||
.await;
|
||||
} else {
|
||||
record_bucket_object_write_memory(&bucket, previous_current_size, obj_info.size.max(0) as u64).await;
|
||||
}
|
||||
}
|
||||
None => {
|
||||
// Neither source could determine the previous state (peers
|
||||
// predating the backfill field during a rolling upgrade, or
|
||||
// sub-quorum metadata divergence). Record the components that
|
||||
// are correct regardless; the next authoritative scanner
|
||||
// refresh replaces the in-memory numbers.
|
||||
debug!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_versioned,
|
||||
"put_object old-size backfill unknown; recording degraded usage delta"
|
||||
);
|
||||
record_bucket_object_write_unknown_previous_memory(&bucket, obj_info.size.max(0) as u64, put_versioned)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
if dsc.replicate_any() {
|
||||
schedule_object_replication(obj_info.clone(), store, dsc).await;
|
||||
}
|
||||
|
||||
rustfs_scanner::record_dirty_usage_bucket(&bucket);
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_post_store_bookkeeping", post_store_stage_start);
|
||||
|
||||
let capacity_update_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
let manager = get_capacity_manager();
|
||||
manager.record_write_operation().await;
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_capacity_update", capacity_update_stage_start);
|
||||
|
||||
Ok::<_, S3Error>(PutObjectCommitResult { obj_info, put_versioned })
|
||||
}
|
||||
});
|
||||
let PutObjectCommitResult { obj_info, put_versioned } = match put_commit.await {
|
||||
Ok(Ok(result)) => result,
|
||||
Ok(Err(err)) => {
|
||||
let result: S3Result<S3Response<PutObjectOutput>> = Err(err);
|
||||
put_request_guard.finish_err();
|
||||
let _ = helper.complete(&result);
|
||||
return result;
|
||||
}
|
||||
Err(err) => {
|
||||
store_put_watchdog.cancel();
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start);
|
||||
warn!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
event = EVENT_PUT_OBJECT_STORE_RETURNED,
|
||||
component = LOG_COMPONENT_APP,
|
||||
subsystem = LOG_SUBSYSTEM_OBJECT,
|
||||
request_id = %request_id,
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_path = put_path,
|
||||
object_size = actual_size,
|
||||
duration_ms = start_time.elapsed().as_millis() as u64,
|
||||
result = "error",
|
||||
error = %err,
|
||||
"PutObject store write returned"
|
||||
);
|
||||
let result: S3Result<S3Response<PutObjectOutput>> = Err(err.into());
|
||||
let result: S3Result<S3Response<PutObjectOutput>> = Err(S3Error::with_message(
|
||||
S3ErrorCode::InternalError,
|
||||
format!("put object commit owner task failed: {err}"),
|
||||
));
|
||||
put_request_guard.finish_err();
|
||||
let _ = helper.complete(&result);
|
||||
return result;
|
||||
}
|
||||
};
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start);
|
||||
drop(object_traffic_progress);
|
||||
#[cfg(test)]
|
||||
wait_for_put_post_store_test_hook(&bucket).await;
|
||||
|
||||
let post_store_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
maybe_enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await;
|
||||
let _ = invalidate_object_data_cache_after_put_success(&cache_adapter, &bucket, &key).await;
|
||||
|
||||
let put_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
|
||||
// Fast in-memory update for immediate quota and admin usage consistency.
|
||||
// The previous current size comes from the prelookup when it ran,
|
||||
// otherwise from the rename_data backfill (rustfs/backlog#1009); the
|
||||
// backfill reproduces the lookup's observation bit for bit (latest
|
||||
// version's ObjectInfo.size — 0 for a delete-marker latest — or
|
||||
// not-found → None).
|
||||
match prelookup_previous_current_size.or_else(|| previous_current_size_from_backfill(backfilled_old_current_size)) {
|
||||
Some(previous_current_size) => {
|
||||
if put_versioned {
|
||||
record_bucket_object_version_write_memory(&bucket, previous_current_size, obj_info.size.max(0) as u64).await;
|
||||
} else {
|
||||
record_bucket_object_write_memory(&bucket, previous_current_size, obj_info.size.max(0) as u64).await;
|
||||
}
|
||||
}
|
||||
None => {
|
||||
// Neither source could determine the previous state (peers
|
||||
// predating the backfill field during a rolling upgrade, or
|
||||
// sub-quorum metadata divergence). Record the components that
|
||||
// are correct regardless; the next authoritative scanner
|
||||
// refresh replaces the in-memory numbers.
|
||||
debug!(
|
||||
target: "rustfs::app::object_usecase",
|
||||
bucket = %bucket,
|
||||
key = %key,
|
||||
put_versioned,
|
||||
"put_object old-size backfill unknown; recording degraded usage delta"
|
||||
);
|
||||
record_bucket_object_write_unknown_previous_memory(&bucket, obj_info.size.max(0) as u64, put_versioned).await;
|
||||
}
|
||||
}
|
||||
|
||||
let raw_version = obj_info.version_id.map(|v| v.to_string());
|
||||
|
||||
@@ -6110,17 +6163,6 @@ impl DefaultObjectUsecase {
|
||||
|
||||
let expiration = resolve_put_object_expiration(&bucket, &obj_info).await;
|
||||
|
||||
// Reuse the single replication decision computed before commit (see `dsc`
|
||||
// above) so the pending metadata persisted with the object and the
|
||||
// post-commit schedule always derive from the same immutable decision.
|
||||
// Recomputing here would repeat the versioning/config/target traversal and,
|
||||
// worse, allow a replication-config hot update between the two phases to
|
||||
// produce a pending-without-schedule or schedule-without-pending divergence
|
||||
// (https://github.com/rustfs/backlog/issues/1320).
|
||||
if dsc.replicate_any() {
|
||||
schedule_object_replication(obj_info.clone(), store, dsc).await;
|
||||
}
|
||||
|
||||
let mut checksums = PutObjectChecksums {
|
||||
crc32: input.checksum_crc32,
|
||||
crc32c: input.checksum_crc32c,
|
||||
@@ -6159,14 +6201,6 @@ impl DefaultObjectUsecase {
|
||||
inject_additional_checksum_headers(&mut response.headers, &put_extra_checksum_headers);
|
||||
let result = Ok(response);
|
||||
let _ = helper.complete(&result);
|
||||
rustfs_scanner::record_dirty_usage_bucket(&bucket);
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_post_store_bookkeeping", post_store_stage_start);
|
||||
|
||||
// Record write operation for capacity management (inline to avoid per-request tokio::spawn overhead)
|
||||
let capacity_update_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
let manager = get_capacity_manager();
|
||||
manager.record_write_operation().await;
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_capacity_update", capacity_update_stage_start);
|
||||
|
||||
// Record PutObject metrics via zero-copy-metrics
|
||||
{
|
||||
@@ -11507,6 +11541,76 @@ mod tests {
|
||||
assert!(!recovered.write_stalled);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(body_cache_hook)]
|
||||
async fn cancelled_put_request_completes_post_commit_publication() {
|
||||
use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions};
|
||||
|
||||
let (store, context) = real_cold_fill_test_context().await;
|
||||
let bucket = format!("put-owner-tail-{}", Uuid::new_v4());
|
||||
let object = "object.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("PUT owner-tail bucket must be created");
|
||||
|
||||
let old_body = Bytes::from_static(b"old body that must be invalidated");
|
||||
let old_info = put_real_cold_fill_object(&store, &bucket, object, &old_body).await;
|
||||
let adapter = context.object_data_cache();
|
||||
let old_plan = real_cold_fill_plan(&adapter, &bucket, object, &old_info);
|
||||
|
||||
let post_store_entered = Arc::new(tokio::sync::Barrier::new(2));
|
||||
let post_store_resume = Arc::new(tokio::sync::Barrier::new(2));
|
||||
install_put_post_store_test_hook(bucket.clone(), Arc::clone(&post_store_entered), Arc::clone(&post_store_resume));
|
||||
|
||||
let payload = Bytes::from_static(b"published despite caller cancellation");
|
||||
let put_input = PutObjectInput::builder()
|
||||
.bucket(bucket.clone())
|
||||
.key(object.to_string())
|
||||
.body(Some(StreamingBlob::from(s3s::Body::from(payload.clone()))))
|
||||
.content_length(Some(i64::try_from(payload.len()).expect("test payload length must fit i64")))
|
||||
.build()
|
||||
.expect("PUT input must build");
|
||||
let put_usecase = DefaultObjectUsecase::with_context(Some(Arc::clone(&context)));
|
||||
let put = tokio::spawn(async move {
|
||||
put_usecase
|
||||
.execute_put_object(&FS::new(), build_request(put_input, Method::PUT))
|
||||
.await
|
||||
});
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(10), post_store_entered.wait())
|
||||
.await
|
||||
.expect("PUT must reach the post-store owner-tail hook");
|
||||
assert_eq!(
|
||||
adapter.fill_body(&old_plan, old_body.clone()).await,
|
||||
rustfs_object_data_cache::ObjectDataCacheFillResult::Inserted,
|
||||
"test must republish the old body while the owner tail is paused"
|
||||
);
|
||||
put.abort();
|
||||
post_store_resume.wait().await;
|
||||
let _ = put.await.expect_err("outer request task must be cancelled");
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(10), async {
|
||||
loop {
|
||||
if matches!(
|
||||
adapter.lookup_body(&old_plan).await,
|
||||
rustfs_object_data_cache::ObjectDataCacheLookup::Miss
|
||||
) {
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("post-commit owner tail must invalidate stale body cache after caller cancellation");
|
||||
|
||||
let recovered = store
|
||||
.get_object_info(&bucket, object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("cancelled request's owned commit must still publish the object");
|
||||
assert_eq!(recovered.size, i64::try_from(payload.len()).expect("test payload length must fit i64"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn object_progress_tracks_zero_byte_and_zero_copy_put_lock_waits() {
|
||||
use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions};
|
||||
|
||||
@@ -1150,7 +1150,9 @@ pub(crate) mod multipart_usecase {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) use super::{access, bucket, data_usage, error, helper, io, object_utils, options, s3_api, set_disk, sse};
|
||||
pub(crate) use super::{
|
||||
access, bucket, data_usage, error, helper, io, object_utils, options, request_context, s3_api, set_disk, sse,
|
||||
};
|
||||
pub(crate) use crate::storage::storage_api::{ECStore, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader};
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user