mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-05 12:57:42 +00:00
fix(tier): recover (#3182)
* fix(tier): stop sending nil/garbage versionId to warm backend S3 Three bugs caused NoSuchVersion errors when reading tiered objects: 1. warm_backend_s3sdk: GET and DELETE ignored rv/range opts entirely — fixed to forward version_id and byte-range to the SDK request. 2. version.rs (MetaObject + MetaDeleteMarker): transition_version_id was parsed with unwrap_or_default(), turning invalid/wrong-length bytes into Uuid::nil(). The nil UUID was then serialized and sent as ?versionId=00000000-... to the tier backend -> NoSuchVersion. Fixed: .and_then(.ok()).filter(!is_nil()) so only valid non-nil UUIDs are forwarded as versionId. 3. bucket_lifecycle_ops: add debug/error logs in get_transitioned_object_reader to record tier, tier_object, and tier_version_id before and on failure of the tier GET. Also adds tier transition fields to dump_fileinfo example for offline xl.meta inspection, and fixes Docker build (cargo path + entrypoint). Adds CLAUDE.md with tier architecture and debugging notes. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * more fixes for versionId * Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: Marcelo Bartsch <marcelo@bartsch.cl> * remove branch * Add tests and fix cargo path, add load to build-docker * update documentation (CLAUDE.md) * more fixes for recover * More fixes to ILM recover * final fix * chore: add missing-shard first-scene diagnostics (#3213) chore(ecstore): add missing-shard first-scene diagnostics Log rename_data quorum context behind RUSTFS_ISSUE3031_DIAG_ENABLE so partial-disk success can be correlated with later missing shard reads. Also log put_object commit success and tmp cleanup boundaries to capture when successful quorum writes are followed by tmp_dir cleanup. * fix test anmd fmt * fix cargo path fix test * fix(tier): format copy_object self-copy guard --------- Signed-off-by: Marcelo Bartsch <marcelo@bartsch.cl> Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com> Co-authored-by: 安正超 <anzhengchao@gmail.com> Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Co-authored-by: houseme <housemecn@gmail.com> Co-authored-by: cxymds <Cxymds@qq.com> Co-authored-by: loverustfs <hello@rustfs.com>
This commit is contained in:
@@ -9,7 +9,7 @@ build-docker: SOURCE_BUILD_CONTAINER_NAME = rustfs-$(BUILD_OS)-build
|
||||
build-docker: BUILD_CMD = /root/.cargo/bin/cargo build --release --bin rustfs --target-dir /root/s3-rustfs/target/$(BUILD_OS)
|
||||
build-docker: ## Build using Docker container # e.g (make build-docker BUILD_OS=ubuntu22.04)
|
||||
@echo "🐳 Building RustFS using Docker ($(BUILD_OS))..."
|
||||
$(DOCKER_CLI) buildx build -t $(SOURCE_BUILD_IMAGE_NAME) -f $(DOCKERFILE_SOURCE) .
|
||||
$(DOCKER_CLI) buildx build -t $(SOURCE_BUILD_IMAGE_NAME) -f $(DOCKERFILE_SOURCE) --load .
|
||||
$(DOCKER_CLI) run --rm --name $(SOURCE_BUILD_CONTAINER_NAME) -v $(shell pwd):/root/s3-rustfs -it $(SOURCE_BUILD_IMAGE_NAME) $(BUILD_CMD)
|
||||
|
||||
.PHONY: docker-inspect-multiarch
|
||||
@@ -19,4 +19,4 @@ docker-inspect-multiarch: ## Check image architecture support
|
||||
exit 1; \
|
||||
fi
|
||||
@echo "🔍 Inspecting multi-architecture image: $(IMAGE)"
|
||||
docker buildx imagetools inspect $(IMAGE)
|
||||
docker buildx imagetools inspect $(IMAGE)
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
# RustFS — CLAUDE.md
|
||||
|
||||
S3-compatible object store in Rust, derived from MinIO. Erasure-coded, multi-pool, supports ILM tiering/lifecycle.
|
||||
|
||||
## Commands
|
||||
|
||||
```bash
|
||||
cargo build --release --bin rustfs # production binary
|
||||
cargo build # dev build
|
||||
cargo check -p <crate> # fast type-check one crate
|
||||
cargo test -p <crate> # test one crate
|
||||
cargo fmt --all # format (required before PR)
|
||||
make pre-commit # full pre-PR gate (fmt + clippy + test)
|
||||
make build-docker BUILD_OS=ubuntu22.04 # Docker cross-build
|
||||
```
|
||||
|
||||
> **Docker build note**: `buildx build` without `--load` keeps the image in the buildx cache only — `docker run` will use a stale local image. The Makefile already includes `--load`; if you suspect a stale binary, add `--no-cache` to the `buildx build` invocation inside `.config/make/build-docker.mak`.
|
||||
|
||||
> Agent/PR rules: see `.github/copilot-instructions.md`.
|
||||
> Crate membership: `Cargo.toml` `[workspace].members`.
|
||||
> CI gates: `.github/workflows/ci.yml`.
|
||||
|
||||
## Workspace layout
|
||||
|
||||
```
|
||||
rustfs/src/main.rs # binary entry point
|
||||
crates/ecstore/src/
|
||||
set_disk.rs # ErasureSet: transition_object, restore_transitioned_object
|
||||
store.rs / store_api/ # ECStore trait + ObjectInfo / TransitionedObject types
|
||||
bucket/lifecycle/
|
||||
bucket_lifecycle_ops.rs # ILM actions: transition_object, expire_transitioned_object,
|
||||
# get_transitioned_object_reader, gen_transition_objname
|
||||
tier_sweeper.rs # background sweep: delete_object_from_remote_tier
|
||||
tier/
|
||||
warm_backend.rs # WarmBackend trait (put/get/remove/in_use)
|
||||
warm_backend_s3.rs # HTTP-client based (TransitionClient) — used for S3/MinIO
|
||||
warm_backend_s3sdk.rs # aws-sdk-s3 based — alternative S3 backend
|
||||
warm_backend_minio.rs / _rustfs.rs / … # per-provider wrappers (all delegate to _s3 or _s3sdk)
|
||||
tier.rs # TierConfigMgr, new_warm_backend dispatch
|
||||
client/transition_api.rs # TransitionClient HTTP plumbing; UploadInfo, to_object_info
|
||||
client/api_put_object_streaming.rs # put_object_do → UploadInfo (version_id from x-amz-version-id)
|
||||
crates/filemeta/src/
|
||||
filemeta.rs # FileMeta (xl.meta top-level), is_skip_meta_key
|
||||
filemeta/version.rs # FileMetaVersion, MetaObject, MetaDeleteMarker
|
||||
# → to_fileinfo() reads transition_version_id
|
||||
# → set_transition() writes raw UUID bytes
|
||||
# → From<FileInfo> for MetaObject writes all meta
|
||||
fileinfo.rs # FileInfo struct (transition_version_id: Option<Uuid>)
|
||||
examples/
|
||||
dump_fileinfo.rs # CLI: parse xl.meta, print transition fields + metadata
|
||||
dump_versions.rs # CLI: list all versions in xl.meta
|
||||
crates/utils/src/http/metadata_compat.rs # SUFFIX_* constants, insert_bytes/get_bytes (dual RustFS+MinIO keys)
|
||||
```
|
||||
|
||||
## Metadata key conventions
|
||||
|
||||
Internal metadata is stored under **both** `x-rustfs-internal-<suffix>` and `x-minio-internal-<suffix>` for MinIO interoperability. `get_bytes` prefers the RustFS key with MinIO fallback.
|
||||
|
||||
Key suffixes (from `metadata_compat.rs`):
|
||||
| Suffix | Meaning |
|
||||
|--------|---------|
|
||||
| `transition-status` | `"complete"` when tiered |
|
||||
| `transitioned-object` | tier key path (without prefix) |
|
||||
| `transitioned-versionID` | S3 version_id returned by tier PUT (16 raw UUID bytes, or absent) |
|
||||
| `transition-tier` | tier name |
|
||||
| `tier-free-versionID` | delete-marker version for free-version sweep |
|
||||
|
||||
## Tier / ILM transition architecture
|
||||
|
||||
### Transition flow (hot → cold)
|
||||
1. `transition_object` (lifecycle_ops) → `ECStore::transition_object` → `set_disk.rs`
|
||||
2. `gen_transition_objname(bucket)` → `{sha256_hash[0..16]}/{uuid[0..2]}/{uuid[2..4]}/{uuid}` (unique per object version)
|
||||
3. `tgt_client.put_with_meta(dest_obj, …)` → returns `rv: String` (remote S3 version_id, or `""`)
|
||||
4. `fi.transition_version_id = if rv.is_empty() { None } else { Some(Uuid::parse_str(&rv)?) }`
|
||||
5. `fi.transitioned_objname = dest_obj` (without tier prefix)
|
||||
6. Written to xl.meta via `MetaObject::from(FileInfo)` → `insert_bytes(SUFFIX_TRANSITIONED_VERSION_ID, uuid.as_bytes())` (16 raw bytes)
|
||||
|
||||
### Tier GET flow (restore/read)
|
||||
`get_transitioned_object_reader` (lifecycle_ops):
|
||||
- reads `oi.transitioned_object.name` (= `fi.transitioned_objname`)
|
||||
- reads `oi.transitioned_object.version_id` (= `fi.transition_version_id.to_string()` or `""`)
|
||||
- calls `warm_backend.get(name, version_id, opts)`
|
||||
- `warm_backend_s3.rs::get`: adds `?versionId=…` only when `rv != ""`
|
||||
|
||||
### Tier prefix handling
|
||||
`WarmBackendS3::get_dest(object)` prepends `self.prefix` to the object name.
|
||||
`transitioned_objname` is stored **without** the prefix — `get_dest` adds it on every call.
|
||||
|
||||
### xl.meta on disk
|
||||
Path: `{disk}/{bucket}/{object}/xl.meta` — one per erasure shard disk.
|
||||
All shards should be identical for a healthy object.
|
||||
|
||||
## Known bugs & fixes
|
||||
|
||||
### Bug 1: `NoSuchVersion` on tier GET — nil UUID sent as versionId
|
||||
**Root cause**: `transitioned-versionID` metadata key exists with empty string value (0 bytes). Old reading code:
|
||||
```rust
|
||||
// OLD — unwrap_or_default() converts 0-byte or wrong-length slice to Uuid::nil()
|
||||
get_bytes(…).map(|v| Uuid::from_slice(v.as_slice()).unwrap_or_default())
|
||||
// → Some(Uuid::nil()) → sends ?versionId=00000000-… → NoSuchVersion
|
||||
```
|
||||
**Fix** (version.rs, `MetaObject::to_fileinfo` + `MetaDeleteMarker::to_fileinfo`):
|
||||
```rust
|
||||
get_bytes(…)
|
||||
.and_then(|v| Uuid::from_slice(v.as_slice()).ok()) // None for wrong-length bytes
|
||||
.filter(|u| !u.is_nil()) // None for nil UUID (old write-back)
|
||||
```
|
||||
**Regression tests** (`crates/filemeta/src/filemeta/version.rs` `mod tests`): 6 tests cover absent key, empty bytes, nil UUID, and valid UUID round-trip for both `MetaObject` and `MetaDeleteMarker` paths.
|
||||
|
||||
### Bug 2: `warm_backend_s3sdk.rs` ignored rv and range opts
|
||||
**Fix**: added `req.version_id(rv)` and `req.range(…)` to GET; `req.version_id(rv)` to DELETE.
|
||||
|
||||
### Bug 4: `set_disk::copy_object` returns 501 for tiered objects (storage class restore)
|
||||
**Root cause**: `set_disk::copy_object` immediately returns `StorageError::NotImplemented` when `src_info.metadata_only = false`. For tiered objects, `metadata_only` is never set to `true` (guarded by `transitioned_object.tier.is_empty()`). So `mc cp --storage-class STANDARD obj obj` on a tiered object always returns 501.
|
||||
**Fix** (`crates/ecstore/src/set_disk.rs`, `copy_object`):
|
||||
```rust
|
||||
if !src_info.metadata_only {
|
||||
if path_join_buf(&[src_bucket, src_object]) == path_join_buf(&[dst_bucket, dst_object]) {
|
||||
if let Some(mut put_reader) = src_info.put_object_reader.take() {
|
||||
return self.put_object(dst_bucket, dst_object, &mut put_reader, dst_opts).await;
|
||||
}
|
||||
}
|
||||
return Err(StorageError::NotImplemented);
|
||||
}
|
||||
```
|
||||
When a self-copy has a `put_object_reader` (data already fetched from tier in `execute_copy_object`), writes it back locally via `put_object`, effectively de-tiering the object.
|
||||
**How `mc cp --storage-class STANDARD` flows**:
|
||||
1. mc sends `PUT /bucket/key` with `x-amz-copy-source`, `x-amz-metadata-directive: REPLACE`, `x-amz-storage-class: STANDARD`
|
||||
2. `execute_copy_object` → `get_object_reader` fetches data from tier backend → stores in `src_info.put_object_reader`
|
||||
3. `store.copy_object(...)` → now calls `put_object` with tier data and STANDARD storage class in `dst_opts`
|
||||
4. New xl.meta written locally with STANDARD class, no tier metadata → object de-tiered
|
||||
|
||||
### Bug 3 (open): race in `expire_transitioned_object`
|
||||
Order is: delete remote tier version → delete local object.
|
||||
A concurrent GET between those two steps fetches a valid stored version_id but the tier version is already gone → `NoSuchVersion`.
|
||||
The proper fix is to delete local metadata first (making the object unreachable) before deleting the remote tier version.
|
||||
|
||||
## Debugging tier issues
|
||||
|
||||
### Inspect xl.meta directly
|
||||
```bash
|
||||
cargo build -p rustfs-filemeta --example dump_fileinfo
|
||||
./target/debug/examples/dump_fileinfo /srv/rustfs/data/disk0/{bucket}/{object}/xl.meta
|
||||
# Shows: transition_status, transition_tier, transitioned_obj, transition_ver_id
|
||||
```
|
||||
`transition_ver_id: <none>` → no versionId will be sent to tier (correct for non-versioned tier bucket).
|
||||
`transition_ver_id: <uuid>` → that UUID will be sent as `?versionId=<uuid>`.
|
||||
|
||||
### Check what versionId is being sent at runtime
|
||||
Enable debug logging:
|
||||
```bash
|
||||
RUST_LOG=rustfs_ecstore::bucket::lifecycle=debug rustfs …
|
||||
```
|
||||
Log line: `fetching transitioned object from tier` (DEBUG before request).
|
||||
Log line: `tier GET failed` (ERROR on failure, includes `tier_version_id`).
|
||||
|
||||
### Metadata key to watch
|
||||
```
|
||||
x-minio-internal-transitioned-versionID= ← empty string = will cause NoSuchVersion with old code
|
||||
x-rustfs-internal-transitioned-versionID= ← same
|
||||
```
|
||||
If both are empty string, the object was transitioned to a non-versioned tier bucket. The versionId should NOT be sent — fixed by Bug 1 above.
|
||||
|
||||
## Common patterns
|
||||
|
||||
### Writing internal metadata (binary values)
|
||||
```rust
|
||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, uuid.as_bytes().to_vec());
|
||||
// stores under both x-rustfs-internal-* and x-minio-internal-* keys
|
||||
```
|
||||
|
||||
### Reading internal metadata (binary values)
|
||||
```rust
|
||||
get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)
|
||||
.and_then(|v| Uuid::from_slice(v.as_slice()).ok())
|
||||
.filter(|u| !u.is_nil())
|
||||
// Returns None for: absent, wrong-length bytes, nil UUID
|
||||
```
|
||||
|
||||
### WarmBackend trait
|
||||
```rust
|
||||
put_with_meta(object, reader, length, meta) -> Result<String> // returns S3 version_id or ""
|
||||
put(object, reader, length) -> Result<String>
|
||||
get(object, rv, opts) -> Result<ReadCloser> // rv="" means no versionId
|
||||
remove(object, rv) -> Result<()>
|
||||
in_use() -> Result<bool>
|
||||
```
|
||||
`rv` = remote version, always pass as empty string when `transition_version_id` is None.
|
||||
@@ -1770,11 +1770,31 @@ pub async fn get_transitioned_object_reader(
|
||||
gopts.length = length;
|
||||
}
|
||||
|
||||
//return Ok(HttpFileReader::new(rs, &oi, opts, &h));
|
||||
//timeTierAction := auditTierActions(oi.transitioned_object.Tier, length)
|
||||
debug!(
|
||||
bucket = %bucket,
|
||||
object = %object,
|
||||
tier = %oi.transitioned_object.tier,
|
||||
tier_object = %oi.transitioned_object.name,
|
||||
tier_version_id = %oi.transitioned_object.version_id,
|
||||
start_offset = gopts.start_offset,
|
||||
length = gopts.length,
|
||||
"fetching transitioned object from tier"
|
||||
);
|
||||
let reader = tgt_client
|
||||
.get(&oi.transitioned_object.name, &oi.transitioned_object.version_id, gopts)
|
||||
.await?;
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!(
|
||||
bucket = %bucket,
|
||||
object = %object,
|
||||
tier = %oi.transitioned_object.tier,
|
||||
tier_object = %oi.transitioned_object.name,
|
||||
tier_version_id = %oi.transitioned_object.version_id,
|
||||
error = %e,
|
||||
"tier GET failed"
|
||||
);
|
||||
e
|
||||
})?;
|
||||
Ok(get_fn(reader, h.clone()))
|
||||
}
|
||||
|
||||
|
||||
@@ -1635,10 +1635,17 @@ impl ObjectOperations for SetDisks {
|
||||
src_opts: &ObjectOptions,
|
||||
dst_opts: &ObjectOptions,
|
||||
) -> Result<ObjectInfo> {
|
||||
// FIXME: TODO:
|
||||
|
||||
if !src_info.metadata_only {
|
||||
return Err(StorageError::NotImplemented);
|
||||
if path_join_buf(&[src_bucket, src_object]) != path_join_buf(&[dst_bucket, dst_object]) {
|
||||
return Err(StorageError::NotImplemented);
|
||||
}
|
||||
// Self-copy with a data reader: write tier data back locally (de-tiering).
|
||||
// Handles `mc cp --storage-class STANDARD obj obj` on a transitioned object.
|
||||
if let Some(mut put_reader) = src_info.put_object_reader.take() {
|
||||
return self.put_object(dst_bucket, dst_object, &mut put_reader, dst_opts).await;
|
||||
}
|
||||
// Same-key tiered copy without a pre-fetched reader: fall through to the metadata
|
||||
// path so the caller gets a disk/quorum error rather than NotImplemented.
|
||||
}
|
||||
|
||||
if path_join_buf(&[src_bucket, src_object]) != path_join_buf(&[dst_bucket, dst_object]) {
|
||||
@@ -4727,6 +4734,7 @@ pub fn is_infrequent_access_class(storage_class: &str) -> bool {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::bucket::lifecycle::bucket_lifecycle_ops::TransitionedObject;
|
||||
use crate::disk::CHECK_PART_UNKNOWN;
|
||||
use crate::disk::CHECK_PART_VOLUME_NOT_FOUND;
|
||||
use crate::disk::RUSTFS_META_BUCKET;
|
||||
@@ -6735,4 +6743,59 @@ mod tests {
|
||||
assert!(!is_infrequent_access_class(storageclass::DEEP_ARCHIVE));
|
||||
assert!(!is_infrequent_access_class(storageclass::EXPRESS_ONEZONE));
|
||||
}
|
||||
|
||||
// Regression test: `mc cp --storage-class STANDARD` on a tiered object (self-copy) must not
|
||||
// return NotImplemented. When the source object is tiered (transitioned_object.tier is
|
||||
// non-empty) the usecase layer in object_usecase.rs intentionally leaves metadata_only=false
|
||||
// so that the full copy path is taken. SetDisks::copy_object must therefore accept a
|
||||
// same-bucket/same-key call even when metadata_only=false.
|
||||
//
|
||||
// Currently this test FAILS because the guard at set_disk.rs:1579 unconditionally rejects
|
||||
// !metadata_only with StorageError::NotImplemented. Once the fix is applied the test will
|
||||
// pass (or progress further through the copy path before failing on missing disk data).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
#[serial]
|
||||
async fn copy_object_tiered_self_copy_does_not_return_not_implemented() {
|
||||
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await;
|
||||
let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new(
|
||||
rustfs_lock::GlobalLockManager::new(),
|
||||
)))])
|
||||
.await;
|
||||
|
||||
// Simulate a tiered object: metadata_only is false (set_disk must handle the full copy),
|
||||
// and transitioned_object.tier is non-empty (the object lives on a remote tier).
|
||||
let mut src_info = ObjectInfo {
|
||||
metadata_only: false,
|
||||
transitioned_object: TransitionedObject {
|
||||
tier: "NEXTCLOUD".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let result = set_disks
|
||||
.copy_object(
|
||||
"bucket",
|
||||
"object",
|
||||
"bucket",
|
||||
"object",
|
||||
&mut src_info,
|
||||
&ObjectOptions::default(),
|
||||
&ObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
// The copy must not be rejected with NotImplemented. Any other outcome (Ok or a
|
||||
// different error such as missing-disk / quorum) is acceptable here.
|
||||
if let Err(ref err) = result {
|
||||
assert!(
|
||||
!matches!(err, StorageError::NotImplemented),
|
||||
"tiered self-copy returned NotImplemented — copy_object must handle \
|
||||
metadata_only=false for same-key copies of tiered objects, got: {err}"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -570,9 +570,32 @@ impl ECStore {
|
||||
}
|
||||
|
||||
if !dst_opts.versioned && src_opts.version_id.is_none() {
|
||||
return self.pools[pool_idx]
|
||||
.copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, &dst_opts)
|
||||
.await;
|
||||
if src_info.metadata_only {
|
||||
return self.pools[pool_idx]
|
||||
.copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, &dst_opts)
|
||||
.await;
|
||||
}
|
||||
// Transitioned object self-copy: restore from tier into the same pool.
|
||||
let put_opts = ObjectOptions {
|
||||
user_defined: (*src_info.user_defined).clone(),
|
||||
versioned: dst_opts.versioned,
|
||||
version_id: dst_opts.version_id.clone(),
|
||||
no_lock: dst_opts.no_lock,
|
||||
mod_time: dst_opts.mod_time,
|
||||
http_preconditions: dst_opts.http_preconditions.clone(),
|
||||
..Default::default()
|
||||
};
|
||||
return if let Some(reader) = src_info.put_object_reader.as_mut() {
|
||||
self.pools[pool_idx]
|
||||
.put_object(dst_bucket, &dst_object, reader, &put_opts)
|
||||
.await
|
||||
} else {
|
||||
Err(StorageError::InvalidArgument(
|
||||
src_bucket.to_owned(),
|
||||
src_object.to_owned(),
|
||||
"put_object_reader is none".to_owned(),
|
||||
))
|
||||
};
|
||||
}
|
||||
|
||||
if dst_opts.versioned && src_opts.version_id != dst_opts.version_id {
|
||||
|
||||
@@ -87,6 +87,7 @@ fn parse_http_timestamp(value: &str) -> Option<OffsetDateTime> {
|
||||
pub fn build_transition_put_options(storage_class: String, mut metadata: HashMap<String, String>) -> PutObjectOptions {
|
||||
let mut opts = PutObjectOptions {
|
||||
storage_class,
|
||||
send_content_md5: true,
|
||||
legalhold: ObjectLockLegalHoldStatus::from_static(""),
|
||||
internal: AdvancedPutOptions {
|
||||
replication_status: ReplicationStatus::from_static(""),
|
||||
|
||||
@@ -147,15 +147,22 @@ impl WarmBackend for WarmBackendS3 {
|
||||
|
||||
async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
|
||||
let client = self.client.clone();
|
||||
let Ok(res) = client
|
||||
.get_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&self.get_dest(object))
|
||||
.send()
|
||||
.await
|
||||
else {
|
||||
return Err(std::io::Error::other("get_object error"));
|
||||
};
|
||||
let mut req = client.get_object().bucket(&self.bucket).key(&self.get_dest(object));
|
||||
|
||||
if !rv.is_empty() {
|
||||
req = req.version_id(rv);
|
||||
}
|
||||
|
||||
if opts.start_offset >= 0 && opts.length > 0 {
|
||||
let end = opts
|
||||
.start_offset
|
||||
.checked_add(opts.length)
|
||||
.and_then(|v| v.checked_sub(1))
|
||||
.ok_or_else(|| std::io::Error::other("invalid range: overflow"))?;
|
||||
req = req.range(format!("bytes={}-{}", opts.start_offset, end));
|
||||
}
|
||||
|
||||
let res = req.send().await.map_err(|e| std::io::Error::other(e.to_string()))?;
|
||||
|
||||
Ok(ReadCloser::new(std::io::Cursor::new(
|
||||
res.body.collect().await.map(|data| data.into_bytes().to_vec())?,
|
||||
@@ -164,16 +171,14 @@ impl WarmBackend for WarmBackendS3 {
|
||||
|
||||
async fn remove(&self, object: &str, rv: &str) -> Result<(), std::io::Error> {
|
||||
let client = self.client.clone();
|
||||
if let Err(_) = client
|
||||
.delete_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&self.get_dest(object))
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
return Err(std::io::Error::other("delete_object error"));
|
||||
let mut req = client.delete_object().bucket(&self.bucket).key(&self.get_dest(object));
|
||||
|
||||
if !rv.is_empty() {
|
||||
req = req.version_id(rv);
|
||||
}
|
||||
|
||||
req.send().await.map_err(|e| std::io::Error::other(e.to_string()))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -41,6 +41,19 @@ fn main() {
|
||||
part.number, part.size, part.actual_size, part.etag
|
||||
);
|
||||
}
|
||||
// Tier / transition fields
|
||||
if !fi.transition_status.is_empty() {
|
||||
println!("transition_status: {}", fi.transition_status);
|
||||
println!("transition_tier: {}", fi.transition_tier);
|
||||
println!("transitioned_obj: {}", fi.transitioned_objname);
|
||||
println!(
|
||||
"transition_ver_id: {}",
|
||||
fi.transition_version_id
|
||||
.map(|u| u.to_string())
|
||||
.unwrap_or_else(|| "<none>".into())
|
||||
);
|
||||
}
|
||||
|
||||
println!("metadata entries: {}", fi.metadata.len());
|
||||
let mut keys = fi.metadata.keys().cloned().collect::<Vec<_>>();
|
||||
keys.sort();
|
||||
|
||||
@@ -2052,8 +2052,9 @@ impl MetaObject {
|
||||
let transitioned_objname = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_OBJECTNAME)
|
||||
.map(|v| String::from_utf8_lossy(&v).to_string())
|
||||
.unwrap_or_default();
|
||||
let transition_version_id =
|
||||
get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID).map(|v| Uuid::from_slice(v.as_slice()).unwrap_or_default());
|
||||
let transition_version_id = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)
|
||||
.and_then(|v| Uuid::from_slice(v.as_slice()).ok())
|
||||
.filter(|u| !u.is_nil());
|
||||
let transition_tier = get_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER)
|
||||
.map(|v| String::from_utf8_lossy(&v).to_string())
|
||||
.unwrap_or_default();
|
||||
@@ -2359,7 +2360,8 @@ impl MetaDeleteMarker {
|
||||
.unwrap_or_default();
|
||||
|
||||
fi.transition_version_id = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)
|
||||
.map(|v| Uuid::from_slice(v.as_slice()).unwrap_or_default());
|
||||
.and_then(|v| Uuid::from_slice(v.as_slice()).ok())
|
||||
.filter(|u| !u.is_nil());
|
||||
}
|
||||
|
||||
fi
|
||||
@@ -3314,4 +3316,79 @@ mod tests {
|
||||
let decoded = FileMetaVersion::decode_data_dir_from_meta(&encoded).expect("decode data_dir");
|
||||
assert_eq!(decoded, Some(data_dir));
|
||||
}
|
||||
|
||||
fn make_meta_object_with_sys(meta_sys: HashMap<String, Vec<u8>>) -> MetaObject {
|
||||
MetaObject {
|
||||
erasure_algorithm: ErasureAlgo::ReedSolomon,
|
||||
erasure_m: 2,
|
||||
erasure_n: 4,
|
||||
erasure_block_size: 1_048_576,
|
||||
erasure_index: 1,
|
||||
erasure_dist: vec![1, 2, 3, 4, 5, 6],
|
||||
bitrot_checksum_algo: ChecksumAlgo::HighwayHash,
|
||||
meta_sys,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn meta_object_transition_version_id_absent_yields_none() {
|
||||
let fi = make_meta_object_with_sys(HashMap::new()).into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn meta_object_transition_version_id_empty_bytes_yields_none() {
|
||||
let mut sys = HashMap::new();
|
||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, vec![]);
|
||||
let fi = make_meta_object_with_sys(sys).into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn meta_object_transition_version_id_nil_uuid_yields_none() {
|
||||
// Regression: old code used unwrap_or_default() which turned nil bytes into Some(Uuid::nil())
|
||||
let mut sys = HashMap::new();
|
||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, Uuid::nil().as_bytes().to_vec());
|
||||
let fi = make_meta_object_with_sys(sys).into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn meta_object_transition_version_id_valid_uuid_round_trips() {
|
||||
let id = sample_version_id();
|
||||
let mut sys = HashMap::new();
|
||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, id.as_bytes().to_vec());
|
||||
let fi = make_meta_object_with_sys(sys).into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, Some(id));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delete_marker_free_version_transition_version_id_nil_uuid_yields_none() {
|
||||
let mut sys = HashMap::new();
|
||||
insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]);
|
||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, Uuid::nil().as_bytes().to_vec());
|
||||
let fi = MetaDeleteMarker {
|
||||
version_id: None,
|
||||
mod_time: None,
|
||||
meta_sys: sys,
|
||||
}
|
||||
.into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delete_marker_free_version_transition_version_id_valid_uuid_round_trips() {
|
||||
let id = sample_version_id();
|
||||
let mut sys = HashMap::new();
|
||||
insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]);
|
||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, id.as_bytes().to_vec());
|
||||
let fi = MetaDeleteMarker {
|
||||
version_id: None,
|
||||
mod_time: None,
|
||||
meta_sys: sys,
|
||||
}
|
||||
.into_fileinfo("b", "k", false);
|
||||
assert_eq!(fi.transition_version_id, Some(id));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2694,6 +2694,7 @@ impl DefaultObjectUsecase {
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
storage_class,
|
||||
..
|
||||
} = req.input.clone();
|
||||
let (src_bucket, src_key, version_id) = match copy_source {
|
||||
@@ -2706,14 +2707,22 @@ impl DefaultObjectUsecase {
|
||||
} => (bucket.to_string(), key.to_string(), version_id.map(|v| v.to_string())),
|
||||
};
|
||||
|
||||
if let Some(ref sc) = storage_class
|
||||
&& !is_valid_storage_class(sc.as_str())
|
||||
{
|
||||
return Err(s3_error!(InvalidStorageClass));
|
||||
}
|
||||
|
||||
// Validate both source and destination keys
|
||||
validate_object_key(&src_key, "COPY (source)")?;
|
||||
validate_object_key(&key, "COPY (dest)")?;
|
||||
validate_table_catalog_object_mutation(&bucket, &key).await?;
|
||||
|
||||
// AWS S3 allows self-copy when metadata directive is REPLACE (used to update metadata in-place).
|
||||
// Reject only when the directive is not REPLACE.
|
||||
// AWS S3 allows self-copy when metadata directive is REPLACE (used to update metadata in-place),
|
||||
// or when an explicit storage class change is requested.
|
||||
// Reject only when neither condition applies.
|
||||
if metadata_directive.as_ref().map(|d| d.as_str()) != Some(MetadataDirective::REPLACE)
|
||||
&& storage_class.is_none()
|
||||
&& src_bucket == bucket
|
||||
&& src_key == key
|
||||
{
|
||||
@@ -2839,7 +2848,7 @@ impl DefaultObjectUsecase {
|
||||
return Err(s3_error!(PreconditionFailed));
|
||||
}
|
||||
|
||||
if cp_src_dst_same {
|
||||
if cp_src_dst_same && src_info.transitioned_object.tier.is_empty() {
|
||||
src_info.metadata_only = true;
|
||||
}
|
||||
|
||||
@@ -2848,6 +2857,11 @@ impl DefaultObjectUsecase {
|
||||
|
||||
strip_managed_encryption_metadata(&mut user_defined);
|
||||
|
||||
if let Some(ref sc) = storage_class {
|
||||
src_info.storage_class = Some(sc.as_str().to_string());
|
||||
user_defined.insert(AMZ_STORAGE_CLASS.to_string(), sc.as_str().to_string());
|
||||
}
|
||||
|
||||
let actual_size = src_info.get_actual_size().map_err(ApiError::from)?;
|
||||
|
||||
let mut length = actual_size;
|
||||
@@ -5257,6 +5271,74 @@ mod tests {
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_copy_object_rejects_invalid_storage_class() {
|
||||
let input = CopyObjectInput::builder()
|
||||
.copy_source(CopySource::Bucket {
|
||||
bucket: "src-bucket".into(),
|
||||
key: "src-key".into(),
|
||||
version_id: None,
|
||||
})
|
||||
.bucket("dst-bucket".to_string())
|
||||
.key("dst-key".to_string())
|
||||
.storage_class(Some(StorageClass::from_static("INVALID")))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let req = build_request(input, Method::PUT);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
|
||||
let err = Box::pin(usecase.execute_copy_object(req)).await.unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_copy_object_allows_self_copy_with_storage_class_change() {
|
||||
let input = CopyObjectInput::builder()
|
||||
.copy_source(CopySource::Bucket {
|
||||
bucket: "test-bucket".into(),
|
||||
key: "test-key".into(),
|
||||
version_id: None,
|
||||
})
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("test-key".to_string())
|
||||
.storage_class(Some(StorageClass::from_static(storageclass::STANDARD_IA)))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let req = build_request(input, Method::PUT);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
|
||||
let err = Box::pin(usecase.execute_copy_object(req)).await.unwrap_err();
|
||||
// Self-copy with explicit storage class change must pass the self-copy guard.
|
||||
assert_ne!(err.code(), &S3ErrorCode::InvalidRequest);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_copy_object_allows_tiered_self_copy_with_storage_class_change() {
|
||||
let input = CopyObjectInput::builder()
|
||||
.copy_source(CopySource::Bucket {
|
||||
bucket: "test-bucket".into(),
|
||||
key: "test-key".into(),
|
||||
version_id: None,
|
||||
})
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("test-key".to_string())
|
||||
.storage_class(Some(StorageClass::from_static(storageclass::STANDARD)))
|
||||
.metadata_directive(Some(MetadataDirective::from_static(MetadataDirective::REPLACE)))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let req = build_request(input, Method::PUT);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
|
||||
let err = Box::pin(usecase.execute_copy_object(req)).await.unwrap_err();
|
||||
// Tiered self-copy with STANDARD storage class must pass all validation checks.
|
||||
// The call fails at store init (no store in unit tests), not at validation.
|
||||
assert_ne!(err.code(), &S3ErrorCode::InvalidRequest);
|
||||
assert_ne!(err.code(), &S3ErrorCode::NotImplemented);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_delete_object_rejects_invalid_object_key() {
|
||||
let input = DeleteObjectInput::builder()
|
||||
|
||||
Reference in New Issue
Block a user