mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-30 10:08:58 +00:00
Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 30c6f8a9ce | |||
| c77c9875c6 | |||
| 7d698abc1f | |||
| a1a65ad65d | |||
| 358af6a8de | |||
| 90d1a15d13 | |||
| 2423ba8e3f | |||
| 294c79c156 | |||
| 3d80578abd |
@@ -10,16 +10,10 @@ never weaken a check to get green.
|
|||||||
|
|
||||||
## `check_layer_dependencies.sh` — layer DAG in `rustfs/src`
|
## `check_layer_dependencies.sh` — layer DAG in `rustfs/src`
|
||||||
|
|
||||||
Enforces `composition (server, startup/init) → interface (admin,
|
Enforces `interface (admin, storage/ecfs, storage/s3_api) → app → infra`; no
|
||||||
storage/ecfs, storage/s3_api) → app → infra`; no upward imports. Server source
|
upward imports. Known legacy violations live in
|
||||||
files are composition roots, while imports of their exported HTTP contracts
|
|
||||||
are classified as interface dependencies. Known legacy violations live in
|
|
||||||
`scripts/layer-dependency-baseline.txt`.
|
`scripts/layer-dependency-baseline.txt`.
|
||||||
|
|
||||||
Dedicated `*_test.rs` and `tests/` modules are outside this production guard.
|
|
||||||
Inline `#[cfg(test)]` imports remain checked under their source file's layer;
|
|
||||||
move architecture-crossing test scaffolding into a dedicated test module.
|
|
||||||
|
|
||||||
- **New violation**: restructure your change so the dependency points
|
- **New violation**: restructure your change so the dependency points
|
||||||
downward (move the shared type/function to the lower layer).
|
downward (move the shared type/function to the lower layer).
|
||||||
- **You legitimately removed a baseline entry**: run
|
- **You legitimately removed a baseline entry**: run
|
||||||
|
|||||||
+56
-12
@@ -156,16 +156,69 @@ jobs:
|
|||||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||||
cache-save-if: ${{ github.ref == 'refs/heads/main' }}
|
cache-save-if: ${{ github.ref == 'refs/heads/main' }}
|
||||||
|
|
||||||
|
- name: Prepare test evidence
|
||||||
|
run: |
|
||||||
|
mkdir -p artifacts/test-and-lint
|
||||||
|
{
|
||||||
|
echo "run_id=${GITHUB_RUN_ID}"
|
||||||
|
echo "job=${GITHUB_JOB}"
|
||||||
|
echo "runner=${RUNNER_NAME}"
|
||||||
|
echo "started_at=$(date --utc --iso-8601=seconds)"
|
||||||
|
} > artifacts/test-and-lint/run-metadata.txt
|
||||||
|
|
||||||
# Clippy runs before the test pass: lint failures are the most common
|
# Clippy runs before the test pass: lint failures are the most common
|
||||||
# CI-only breakage and should surface in minutes, not after 20+ minutes
|
# CI-only breakage and should surface in minutes, not after 20+ minutes
|
||||||
# of tests.
|
# of tests.
|
||||||
- name: Run clippy lints
|
- name: Run clippy lints
|
||||||
run: cargo clippy --all-targets -- -D warnings
|
run: cargo clippy --all-targets -- -D warnings
|
||||||
|
|
||||||
- name: Run tests
|
- name: Run nextest tests
|
||||||
run: |
|
run: |
|
||||||
cargo nextest run --profile ci --all --exclude e2e_test
|
mkdir -p artifacts/test-and-lint
|
||||||
cargo test --all --doc
|
set +e
|
||||||
|
NEXTEST_HIDE_PROGRESS_BAR=1 timeout --verbose --signal=TERM --kill-after=30s 75m \
|
||||||
|
cargo nextest run --profile ci --all --exclude e2e_test \
|
||||||
|
--status-level all --final-status-level all \
|
||||||
|
2>&1 | tee artifacts/test-and-lint/nextest.log
|
||||||
|
status=${PIPESTATUS[0]}
|
||||||
|
{
|
||||||
|
echo "command=cargo nextest run --profile ci --all --exclude e2e_test"
|
||||||
|
echo "exit_status=${status}"
|
||||||
|
echo "finished_at=$(date --utc --iso-8601=seconds)"
|
||||||
|
echo
|
||||||
|
echo "Remaining test-related processes:"
|
||||||
|
pgrep -af 'cargo|nextest|target/.*/deps/' || true
|
||||||
|
} > artifacts/test-and-lint/nextest-diagnostics.txt
|
||||||
|
exit "${status}"
|
||||||
|
|
||||||
|
- name: Run documentation tests
|
||||||
|
run: |
|
||||||
|
mkdir -p artifacts/test-and-lint
|
||||||
|
set +e
|
||||||
|
timeout --verbose --signal=TERM --kill-after=30s 15m \
|
||||||
|
cargo test --all --doc \
|
||||||
|
2>&1 | tee artifacts/test-and-lint/doctest.log
|
||||||
|
status=${PIPESTATUS[0]}
|
||||||
|
{
|
||||||
|
echo "command=cargo test --all --doc"
|
||||||
|
echo "exit_status=${status}"
|
||||||
|
echo "finished_at=$(date --utc --iso-8601=seconds)"
|
||||||
|
echo
|
||||||
|
echo "Remaining test-related processes:"
|
||||||
|
pgrep -af 'cargo|rustdoc|target/.*/deps/' || true
|
||||||
|
} > artifacts/test-and-lint/doctest-diagnostics.txt
|
||||||
|
exit "${status}"
|
||||||
|
|
||||||
|
- name: Upload test reports and diagnostics
|
||||||
|
if: always()
|
||||||
|
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
|
||||||
|
with:
|
||||||
|
name: junit-test-and-lint-${{ github.run_number }}
|
||||||
|
path: |
|
||||||
|
target/nextest/ci/junit.xml
|
||||||
|
artifacts/test-and-lint
|
||||||
|
retention-days: 3
|
||||||
|
if-no-files-found: error
|
||||||
|
|
||||||
# rustfs/backlog#1289: fail if a seed rule's log anchor no longer exists
|
# rustfs/backlog#1289: fail if a seed rule's log anchor no longer exists
|
||||||
# verbatim in the source tree (log message drifted without updating the
|
# verbatim in the source tree (log message drifted without updating the
|
||||||
@@ -174,15 +227,6 @@ jobs:
|
|||||||
- name: Check log-analyzer rule anchors
|
- name: Check log-analyzer rule anchors
|
||||||
run: ./scripts/check_log_analyzer_rules.sh
|
run: ./scripts/check_log_analyzer_rules.sh
|
||||||
|
|
||||||
- name: Upload test junit report
|
|
||||||
if: always()
|
|
||||||
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
|
|
||||||
with:
|
|
||||||
name: junit-test-and-lint-${{ github.run_number }}
|
|
||||||
path: target/nextest/ci/junit.xml
|
|
||||||
retention-days: 3
|
|
||||||
if-no-files-found: ignore
|
|
||||||
|
|
||||||
# Explicit gate for migration-critical suites. These tests already ran in
|
# Explicit gate for migration-critical suites. These tests already ran in
|
||||||
# the full nextest pass above; a single filtered nextest invocation keeps
|
# the full nextest pass above; a single filtered nextest invocation keeps
|
||||||
# the named gate without rebuilding or re-running them one package at a time.
|
# the named gate without rebuilding or re-running them one package at a time.
|
||||||
|
|||||||
@@ -66,7 +66,7 @@ env:
|
|||||||
CARGO_TERM_COLOR: always
|
CARGO_TERM_COLOR: always
|
||||||
REGISTRY_DOCKERHUB: rustfs/rustfs
|
REGISTRY_DOCKERHUB: rustfs/rustfs
|
||||||
REGISTRY_GHCR: ghcr.io/${{ github.repository }}
|
REGISTRY_GHCR: ghcr.io/${{ github.repository }}
|
||||||
REGISTRY_QUAY: quay.io/${{ secrets.QUAY_USERNAME }}/rustfs
|
REGISTRY_QUAY: quay.io/rustfs/rustfs
|
||||||
DOCKER_PLATFORMS: linux/amd64,linux/arm64
|
DOCKER_PLATFORMS: linux/amd64,linux/arm64
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
|
|||||||
@@ -25,6 +25,9 @@ keywords = ["checksum-calculation", "verification", "integrity", "authenticity",
|
|||||||
categories = ["web-programming", "development-tools", "network-programming"]
|
categories = ["web-programming", "development-tools", "network-programming"]
|
||||||
documentation = "https://docs.rs/rustfs-checksums/latest/rustfs_checksum/"
|
documentation = "https://docs.rs/rustfs-checksums/latest/rustfs_checksum/"
|
||||||
|
|
||||||
|
[lints]
|
||||||
|
workspace = true
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
bytes = { workspace = true, features = ["serde"] }
|
bytes = { workspace = true, features = ["serde"] }
|
||||||
crc-fast = { workspace = true }
|
crc-fast = { workspace = true }
|
||||||
|
|||||||
@@ -506,8 +506,35 @@ pub struct ReplicationStats {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl ReplicationStats {
|
impl ReplicationStats {
|
||||||
|
pub fn is_empty(&self) -> bool {
|
||||||
|
let Self {
|
||||||
|
pending_size,
|
||||||
|
replicated_size,
|
||||||
|
failed_size,
|
||||||
|
failed_count,
|
||||||
|
pending_count,
|
||||||
|
missed_threshold_size,
|
||||||
|
after_threshold_size,
|
||||||
|
missed_threshold_count,
|
||||||
|
after_threshold_count,
|
||||||
|
replicated_count,
|
||||||
|
} = self;
|
||||||
|
|
||||||
|
*pending_size == 0
|
||||||
|
&& *replicated_size == 0
|
||||||
|
&& *failed_size == 0
|
||||||
|
&& *failed_count == 0
|
||||||
|
&& *pending_count == 0
|
||||||
|
&& *missed_threshold_size == 0
|
||||||
|
&& *after_threshold_size == 0
|
||||||
|
&& *missed_threshold_count == 0
|
||||||
|
&& *after_threshold_count == 0
|
||||||
|
&& *replicated_count == 0
|
||||||
|
}
|
||||||
|
|
||||||
|
#[deprecated(note = "use is_empty instead")]
|
||||||
pub fn empty(&self) -> bool {
|
pub fn empty(&self) -> bool {
|
||||||
self.replicated_size == 0 && self.failed_size == 0 && self.failed_count == 0
|
self.is_empty()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -520,16 +547,19 @@ pub struct ReplicationAllStats {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl ReplicationAllStats {
|
impl ReplicationAllStats {
|
||||||
|
pub fn is_empty(&self) -> bool {
|
||||||
|
let Self {
|
||||||
|
replica_size,
|
||||||
|
replica_count,
|
||||||
|
targets,
|
||||||
|
} = self;
|
||||||
|
|
||||||
|
*replica_size == 0 && *replica_count == 0 && targets.values().all(ReplicationStats::is_empty)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[deprecated(note = "use is_empty instead")]
|
||||||
pub fn empty(&self) -> bool {
|
pub fn empty(&self) -> bool {
|
||||||
if self.replica_size != 0 && self.replica_count != 0 {
|
self.is_empty()
|
||||||
return false;
|
|
||||||
}
|
|
||||||
for v in self.targets.values() {
|
|
||||||
if !v.empty() {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
true
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -783,7 +813,7 @@ impl DataUsageCache {
|
|||||||
return Some(root);
|
return Some(root);
|
||||||
}
|
}
|
||||||
let mut flat = self.flatten(&root);
|
let mut flat = self.flatten(&root);
|
||||||
if flat.replication_stats.as_ref().is_some_and(|stats| stats.empty()) {
|
if flat.replication_stats.as_ref().is_some_and(ReplicationAllStats::is_empty) {
|
||||||
flat.replication_stats = None;
|
flat.replication_stats = None;
|
||||||
}
|
}
|
||||||
Some(flat)
|
Some(flat)
|
||||||
@@ -1582,6 +1612,126 @@ mod tests {
|
|||||||
assert_eq!(map["BETWEEN_1024B_AND_1_MB"], u64::MAX);
|
assert_eq!(map["BETWEEN_1024B_AND_1_MB"], u64::MAX);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn replication_stats_empty_checks_every_field() {
|
||||||
|
type SetField = fn(&mut ReplicationStats);
|
||||||
|
|
||||||
|
let cases: [(&str, SetField); 10] = [
|
||||||
|
("pending_size", |stats| stats.pending_size = 1),
|
||||||
|
("replicated_size", |stats| stats.replicated_size = 1),
|
||||||
|
("failed_size", |stats| stats.failed_size = 1),
|
||||||
|
("failed_count", |stats| stats.failed_count = 1),
|
||||||
|
("pending_count", |stats| stats.pending_count = 1),
|
||||||
|
("missed_threshold_size", |stats| stats.missed_threshold_size = 1),
|
||||||
|
("after_threshold_size", |stats| stats.after_threshold_size = 1),
|
||||||
|
("missed_threshold_count", |stats| stats.missed_threshold_count = 1),
|
||||||
|
("after_threshold_count", |stats| stats.after_threshold_count = 1),
|
||||||
|
("replicated_count", |stats| stats.replicated_count = 1),
|
||||||
|
];
|
||||||
|
|
||||||
|
assert!(ReplicationStats::default().is_empty());
|
||||||
|
for (field, set_nonzero) in cases {
|
||||||
|
let mut stats = ReplicationStats::default();
|
||||||
|
set_nonzero(&mut stats);
|
||||||
|
assert!(!stats.is_empty(), "{field} must make replication stats non-empty");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn replication_all_stats_empty_checks_aggregate_fields_independently() {
|
||||||
|
let cases = [
|
||||||
|
(
|
||||||
|
"replica_size",
|
||||||
|
ReplicationAllStats {
|
||||||
|
replica_size: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"replica_count",
|
||||||
|
ReplicationAllStats {
|
||||||
|
replica_count: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
),
|
||||||
|
];
|
||||||
|
|
||||||
|
assert!(ReplicationAllStats::default().is_empty());
|
||||||
|
for (field, stats) in cases {
|
||||||
|
assert!(!stats.is_empty(), "{field} must make aggregate replication stats non-empty");
|
||||||
|
}
|
||||||
|
|
||||||
|
let empty_targets = ReplicationAllStats {
|
||||||
|
targets: HashMap::from([("arn:test:empty".to_string(), ReplicationStats::default())]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(empty_targets.is_empty(), "all-empty targets must keep aggregate stats empty");
|
||||||
|
|
||||||
|
let stats = ReplicationAllStats {
|
||||||
|
targets: HashMap::from([
|
||||||
|
("arn:test:empty".to_string(), ReplicationStats::default()),
|
||||||
|
(
|
||||||
|
"arn:test:non-empty".to_string(),
|
||||||
|
ReplicationStats {
|
||||||
|
pending_count: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
),
|
||||||
|
]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(!stats.is_empty(), "a non-empty target must make aggregate replication stats non-empty");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn size_recursive_prunes_empty_and_preserves_pending_replication_stats() {
|
||||||
|
let root = hash_path("bucket");
|
||||||
|
let child = hash_path("bucket/child");
|
||||||
|
let mut cache = DataUsageCache::default();
|
||||||
|
cache.replace_hashed(&root, &None, &DataUsageEntry::default());
|
||||||
|
cache.replace_hashed(
|
||||||
|
&child,
|
||||||
|
&Some(root.clone()),
|
||||||
|
&DataUsageEntry {
|
||||||
|
replication_stats: Some(ReplicationAllStats::default()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
cache
|
||||||
|
.size_recursive("bucket")
|
||||||
|
.expect("bucket usage should flatten")
|
||||||
|
.replication_stats
|
||||||
|
.is_none()
|
||||||
|
);
|
||||||
|
|
||||||
|
cache.replace_hashed(
|
||||||
|
&child,
|
||||||
|
&Some(root.clone()),
|
||||||
|
&DataUsageEntry {
|
||||||
|
replication_stats: Some(ReplicationAllStats {
|
||||||
|
targets: HashMap::from([(
|
||||||
|
"arn:test:pending".to_string(),
|
||||||
|
ReplicationStats {
|
||||||
|
pending_count: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)]),
|
||||||
|
..Default::default()
|
||||||
|
}),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let flattened = cache.size_recursive("bucket").expect("bucket usage should flatten");
|
||||||
|
let replication = flattened
|
||||||
|
.replication_stats
|
||||||
|
.expect("pending-only replication stats must survive pruning");
|
||||||
|
|
||||||
|
assert_eq!(replication.targets["arn:test:pending"].pending_count, 1);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_data_usage_cache_merge_adds_missing_child() {
|
fn test_data_usage_cache_merge_adds_missing_child() {
|
||||||
let mut base = DataUsageCache::default();
|
let mut base = DataUsageCache::default();
|
||||||
|
|||||||
@@ -554,6 +554,7 @@ async fn delete_free_version_remote_object(
|
|||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
||||||
) -> Result<(), std::io::Error> {
|
) -> Result<(), std::io::Error> {
|
||||||
|
let version_id_exact = validate_transition_remote_version(oi)?;
|
||||||
let identity = tier_destination_id_from_metadata(&oi.user_defined)?
|
let identity = tier_destination_id_from_metadata(&oi.user_defined)?
|
||||||
.ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?;
|
.ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?;
|
||||||
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
||||||
@@ -562,7 +563,7 @@ async fn delete_free_version_remote_object(
|
|||||||
&oi.transitioned_object.tier,
|
&oi.transitioned_object.tier,
|
||||||
identity,
|
identity,
|
||||||
tier_config_mgr,
|
tier_config_mgr,
|
||||||
false,
|
version_id_exact,
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -4201,6 +4202,23 @@ pub async fn get_transitioned_object_reader(
|
|||||||
get_transitioned_object_reader_with_tier_manager(bucket, object, rs, h, oi, opts, &tier_config_mgr).await
|
get_transitioned_object_reader_with_tier_manager(bucket, object, rs, h, oi, opts, &tier_config_mgr).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn validate_transition_remote_version(oi: &ObjectInfo) -> Result<bool, std::io::Error> {
|
||||||
|
let version = oi.transitioned_object.version_id.as_str();
|
||||||
|
match oi.transition_version_state {
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown => Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object version state is unknown",
|
||||||
|
)),
|
||||||
|
rustfs_filemeta::TransitionVersionState::KnownDisabled if version.is_empty() => Ok(false),
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull if version == "null" => Ok(true),
|
||||||
|
rustfs_filemeta::TransitionVersionState::Exact if !version.is_empty() && version != "null" => Ok(true),
|
||||||
|
_ => Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier object version state conflicts with its version ID",
|
||||||
|
)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
@@ -4210,6 +4228,7 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
|
|||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
|
||||||
) -> Result<GetObjectReader, std::io::Error> {
|
) -> Result<GetObjectReader, std::io::Error> {
|
||||||
|
validate_transition_remote_version(oi)?;
|
||||||
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
|
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
|
||||||
let lease = match expected_identity {
|
let lease = match expected_identity {
|
||||||
Some(identity) => {
|
Some(identity) => {
|
||||||
@@ -5506,6 +5525,7 @@ mod tests {
|
|||||||
tier: tier.clone(),
|
tier: tier.clone(),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -5569,6 +5589,7 @@ mod tests {
|
|||||||
tier,
|
tier,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -5591,6 +5612,70 @@ mod tests {
|
|||||||
assert_eq!(backend.get_count().await, 0);
|
assert_eq!(backend.get_count().await, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn transitioned_get_rejects_unknown_version_state_before_backend_io() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&manager, &tier).await;
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
size: 1,
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: "remote/object".to_string(),
|
||||||
|
version_id: String::new(),
|
||||||
|
status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
tier,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = match get_transitioned_object_reader_with_tier_manager(
|
||||||
|
&object_info.bucket,
|
||||||
|
&object_info.name,
|
||||||
|
&None,
|
||||||
|
&HeaderMap::new(),
|
||||||
|
&object_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&manager,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(_) => panic!("unknown remote version state must fail before backend IO"),
|
||||||
|
Err(err) => err,
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(backend.get_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn free_version_delete_rejects_unknown_version_state_before_backend_io() {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let backend = register_mock_tier(&manager, "WARM").await;
|
||||||
|
let object_info = ObjectInfo {
|
||||||
|
transitioned_object: TransitionedObject {
|
||||||
|
name: "remote/object".to_string(),
|
||||||
|
version_id: "legacy-version".to_string(),
|
||||||
|
tier: "WARM".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = super::delete_free_version_remote_object(&object_info, &manager)
|
||||||
|
.await
|
||||||
|
.expect_err("unknown remote version state must fail before backend IO");
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn free_version_remote_delete_requires_persisted_destination_identity() {
|
async fn free_version_remote_delete_requires_persisted_destination_identity() {
|
||||||
@@ -5640,6 +5725,7 @@ mod tests {
|
|||||||
oi.transitioned_object.tier = "WARM".to_string();
|
oi.transitioned_object.tier = "WARM".to_string();
|
||||||
oi.transitioned_object.name = "remote/object".to_string();
|
oi.transitioned_object.name = "remote/object".to_string();
|
||||||
oi.transitioned_object.version_id = "remote-version".to_string();
|
oi.transitioned_object.version_id = "remote-version".to_string();
|
||||||
|
oi.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
|
||||||
let local_delete_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
let local_delete_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||||
|
|
||||||
let legacy_err = delete_free_version_remote_object_then(&oi, &manager, {
|
let legacy_err = delete_free_version_remote_object_then(&oi, &manager, {
|
||||||
@@ -5650,7 +5736,8 @@ mod tests {
|
|||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
.expect_err("legacy free-version without identity must be retained");
|
.expect_err("legacy free-version without identity must be retained");
|
||||||
assert!(legacy_err.to_string().contains("no durable backend identity"));
|
assert_eq!(legacy_err.kind(), std::io::ErrorKind::Other);
|
||||||
|
assert_eq!(old_backend.remove_count().await, 0);
|
||||||
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 0);
|
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 0);
|
||||||
|
|
||||||
let mut invalid_metadata = HashMap::new();
|
let mut invalid_metadata = HashMap::new();
|
||||||
@@ -5781,6 +5868,7 @@ mod tests {
|
|||||||
oi.transitioned_object.tier = "WARM".to_string();
|
oi.transitioned_object.tier = "WARM".to_string();
|
||||||
oi.transitioned_object.name = "remote/object".to_string();
|
oi.transitioned_object.name = "remote/object".to_string();
|
||||||
oi.transitioned_object.version_id = "remote-version".to_string();
|
oi.transitioned_object.version_id = "remote-version".to_string();
|
||||||
|
oi.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
|
||||||
|
|
||||||
let err = match get_transitioned_object_reader_with_tier_manager(
|
let err = match get_transitioned_object_reader_with_tier_manager(
|
||||||
"bucket",
|
"bucket",
|
||||||
@@ -5796,7 +5884,12 @@ mod tests {
|
|||||||
Ok(_) => panic!("identity-bound GET must reject a same-name tier rebind"),
|
Ok(_) => panic!("identity-bound GET must reject a same-name tier rebind"),
|
||||||
Err(err) => err,
|
Err(err) => err,
|
||||||
};
|
};
|
||||||
assert!(err.to_string().contains("identity no longer matches"));
|
assert_eq!(err.kind(), std::io::ErrorKind::Other);
|
||||||
|
let admin_err = err
|
||||||
|
.get_ref()
|
||||||
|
.and_then(|source| source.downcast_ref::<crate::client::admin_handler_utils::AdminError>())
|
||||||
|
.expect("identity mismatch should retain the typed tier error");
|
||||||
|
assert_eq!(admin_err.code, crate::services::tier::tier::ERR_TIER_INVALID_CONFIG.code);
|
||||||
assert_eq!(new_backend.get_count().await, 0);
|
assert_eq!(new_backend.get_count().await, 0);
|
||||||
|
|
||||||
oi.user_defined = Arc::new(HashMap::new());
|
oi.user_defined = Arc::new(HashMap::new());
|
||||||
@@ -5902,7 +5995,8 @@ mod tests {
|
|||||||
version_id: "remote-version".to_string(),
|
version_id: "remote-version".to_string(),
|
||||||
tier_name: "WARM".to_string(),
|
tier_name: "WARM".to_string(),
|
||||||
backend_identity: Some([1; 32]),
|
backend_identity: Some([1; 32]),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
};
|
};
|
||||||
|
|
||||||
let err = state
|
let err = state
|
||||||
@@ -6013,7 +6107,8 @@ mod tests {
|
|||||||
version_id: "remote-version".to_string(),
|
version_id: "remote-version".to_string(),
|
||||||
tier_name: "WARM".to_string(),
|
tier_name: "WARM".to_string(),
|
||||||
backend_identity: Some([1; 32]),
|
backend_identity: Some([1; 32]),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
};
|
};
|
||||||
|
|
||||||
state
|
state
|
||||||
@@ -10081,6 +10176,87 @@ mod tests {
|
|||||||
(backend, identity_hex)
|
(backend, identity_hex)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn journal_replay_rejects_unknown_version_state_before_backend_io() {
|
||||||
|
let (_disk_paths, ecstore) = setup_test_env().await;
|
||||||
|
let (backend, _) = register_recovery_mock_tier(&ecstore).await;
|
||||||
|
let identity = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM")
|
||||||
|
.await
|
||||||
|
.expect("mock tier lease should be available")
|
||||||
|
.backend_identity();
|
||||||
|
let je = Jentry {
|
||||||
|
obj_name: "remote/object".to_string(),
|
||||||
|
version_id: "legacy-version".to_string(),
|
||||||
|
tier_name: "WARM".to_string(),
|
||||||
|
backend_identity: Some(identity),
|
||||||
|
version_id_exact: false,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = crate::bucket::lifecycle::tier_delete_journal::process_tier_delete_journal_entry(ecstore, &je)
|
||||||
|
.await
|
||||||
|
.expect_err("unknown journal state must fail before backend IO");
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
async fn journal_replay_deletes_confirmed_exact_provider_token() {
|
||||||
|
let (_disk_paths, ecstore) = setup_test_env().await;
|
||||||
|
let (backend, _) = register_recovery_mock_tier(&ecstore).await;
|
||||||
|
let lease = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM")
|
||||||
|
.await
|
||||||
|
.expect("mock tier lease should be available");
|
||||||
|
let identity = lease.backend_identity();
|
||||||
|
backend
|
||||||
|
.set_put_remote_version(Some("provider-version-token".to_string()))
|
||||||
|
.await;
|
||||||
|
lease
|
||||||
|
.put(
|
||||||
|
"remote/object",
|
||||||
|
crate::client::transition_api::ReaderImpl::Body(bytes::Bytes::from_static(b"candidate")),
|
||||||
|
9,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("confirmed remote candidate should be seeded");
|
||||||
|
backend.set_remove_failure(true);
|
||||||
|
backend.set_reject_non_empty_remote_versions(true);
|
||||||
|
let je = Jentry {
|
||||||
|
obj_name: "remote/object".to_string(),
|
||||||
|
version_id: "provider-version-token".to_string(),
|
||||||
|
tier_name: "WARM".to_string(),
|
||||||
|
backend_identity: Some(identity),
|
||||||
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
|
};
|
||||||
|
|
||||||
|
crate::set_disk::cleanup_rejected_transition_upload_durably(
|
||||||
|
&lease,
|
||||||
|
&je.obj_name,
|
||||||
|
&je.version_id,
|
||||||
|
true,
|
||||||
|
Some(ecstore.clone()),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("failed immediate cleanup should remain durable in the journal");
|
||||||
|
assert!(backend.contains(&je.obj_name).await);
|
||||||
|
|
||||||
|
backend.set_remove_failure(false);
|
||||||
|
crate::bucket::lifecycle::tier_delete_journal::process_tier_delete_journal_entry(ecstore, &je)
|
||||||
|
.await
|
||||||
|
.expect("identity-bound exact journal must retry confirmed candidate cleanup");
|
||||||
|
|
||||||
|
assert!(!backend.contains(&je.obj_name).await);
|
||||||
|
assert_eq!(backend.exact_remove_count(), 2);
|
||||||
|
assert_eq!(
|
||||||
|
backend.remove_versions().await,
|
||||||
|
vec![("remote/object".to_string(), "provider-version-token".to_string())]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
async fn seed_recoverable_free_version(
|
async fn seed_recoverable_free_version(
|
||||||
disk_paths: &[PathBuf],
|
disk_paths: &[PathBuf],
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -10100,6 +10276,7 @@ mod tests {
|
|||||||
identity,
|
identity,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
let transition_version_id = Uuid::new_v4();
|
||||||
let mut metadata = FileMeta::new();
|
let mut metadata = FileMeta::new();
|
||||||
metadata
|
metadata
|
||||||
.add_version(FileInfo {
|
.add_version(FileInfo {
|
||||||
@@ -10108,7 +10285,9 @@ mod tests {
|
|||||||
version_id: Some(object_version_id),
|
version_id: Some(object_version_id),
|
||||||
transition_status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
transition_status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
transitioned_objname: format!("remote/{bucket}/{object}"),
|
transitioned_objname: format!("remote/{bucket}/{object}"),
|
||||||
transition_version_id: Some(Uuid::new_v4()),
|
transition_version_id: Some(transition_version_id),
|
||||||
|
transition_version: Some(transition_version_id.to_string()),
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
transition_tier: "WARM".to_string(),
|
transition_tier: "WARM".to_string(),
|
||||||
mod_time: Some(OffsetDateTime::now_utc()),
|
mod_time: Some(OffsetDateTime::now_utc()),
|
||||||
metadata: transitioned_metadata,
|
metadata: transitioned_metadata,
|
||||||
|
|||||||
@@ -20,7 +20,10 @@ use tokio_util::sync::CancellationToken;
|
|||||||
use tracing::{debug, warn};
|
use tracing::{debug, warn};
|
||||||
|
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent_with_manager_and_identity};
|
use crate::bucket::lifecycle::tier_sweeper::{
|
||||||
|
Jentry, delete_confirmed_transition_candidate_exact_with_manager_and_identity,
|
||||||
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
||||||
|
};
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
use crate::disk::RUSTFS_META_BUCKET;
|
||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
|
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
|
||||||
@@ -42,6 +45,7 @@ const TIER_DELETE_JOURNAL_RECOVERY_INTERVAL: Duration = Duration::from_secs(60);
|
|||||||
const TIER_DELETE_JOURNAL_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
const TIER_DELETE_JOURNAL_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
||||||
const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
|
const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
|
||||||
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
|
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
|
||||||
|
const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4;
|
||||||
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
|
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||||
@@ -55,24 +59,35 @@ struct PersistedTierDeleteJournalEntry {
|
|||||||
backend_identity: Option<[u8; 32]>,
|
backend_identity: Option<[u8; 32]>,
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
version_id_exact: Option<bool>,
|
version_id_exact: Option<bool>,
|
||||||
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
|
version_state: Option<rustfs_filemeta::TransitionVersionState>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl PersistedTierDeleteJournalEntry {
|
impl PersistedTierDeleteJournalEntry {
|
||||||
fn from_jentry(je: &Jentry) -> Self {
|
fn from_jentry(je: &Jentry) -> Result<Self> {
|
||||||
Self {
|
validate_version_state(je.version_state, &je.version_id, je.version_id_exact)?;
|
||||||
version: if je.version_id_exact {
|
let legacy_unknown = je.version_state == rustfs_filemeta::TransitionVersionState::Unknown;
|
||||||
TIER_DELETE_JOURNAL_EXACT_VERSION
|
let version = if legacy_unknown {
|
||||||
} else if je.backend_identity.is_some() {
|
if je.backend_identity.is_some() {
|
||||||
TIER_DELETE_JOURNAL_VERSION
|
TIER_DELETE_JOURNAL_VERSION
|
||||||
} else {
|
} else {
|
||||||
1
|
1
|
||||||
},
|
}
|
||||||
|
} else {
|
||||||
|
if je.backend_identity.is_none() {
|
||||||
|
return Err(Error::other("new tier delete journal entry is missing its backend identity"));
|
||||||
|
}
|
||||||
|
TIER_DELETE_JOURNAL_STATE_VERSION
|
||||||
|
};
|
||||||
|
Ok(Self {
|
||||||
|
version,
|
||||||
obj_name: je.obj_name.clone(),
|
obj_name: je.obj_name.clone(),
|
||||||
version_id: je.version_id.clone(),
|
version_id: je.version_id.clone(),
|
||||||
tier_name: je.tier_name.clone(),
|
tier_name: je.tier_name.clone(),
|
||||||
backend_identity: je.backend_identity,
|
backend_identity: je.backend_identity,
|
||||||
version_id_exact: je.version_id_exact.then_some(true),
|
version_id_exact: je.version_id_exact.then_some(true),
|
||||||
}
|
version_state: (!legacy_unknown).then_some(je.version_state),
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
fn into_jentry(self) -> Result<Jentry> {
|
fn into_jentry(self) -> Result<Jentry> {
|
||||||
@@ -84,19 +99,23 @@ impl PersistedTierDeleteJournalEntry {
|
|||||||
if self.obj_name.is_empty() || self.tier_name.is_empty() {
|
if self.obj_name.is_empty() || self.tier_name.is_empty() {
|
||||||
return Err(Error::other("tier delete journal entry is incomplete"));
|
return Err(Error::other("tier delete journal entry is incomplete"));
|
||||||
}
|
}
|
||||||
if self.version != TIER_DELETE_JOURNAL_EXACT_VERSION && self.version_id_exact.unwrap_or(false) {
|
if self.version != TIER_DELETE_JOURNAL_EXACT_VERSION
|
||||||
|
&& self.version != TIER_DELETE_JOURNAL_STATE_VERSION
|
||||||
|
&& self.version_id_exact.unwrap_or(false)
|
||||||
|
{
|
||||||
return Err(Error::other(
|
return Err(Error::other(
|
||||||
"legacy tier delete journal entry has an unsupported exact version constraint",
|
"legacy tier delete journal entry has an unsupported exact version constraint",
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
let (backend_identity, version_id_exact) = match self.version {
|
let (backend_identity, version_id_exact, version_state) = match self.version {
|
||||||
1 => (None, false),
|
1 => (None, false, rustfs_filemeta::TransitionVersionState::Unknown),
|
||||||
TIER_DELETE_JOURNAL_VERSION => (
|
TIER_DELETE_JOURNAL_VERSION => (
|
||||||
Some(
|
Some(
|
||||||
self.backend_identity
|
self.backend_identity
|
||||||
.ok_or_else(|| Error::other("tier delete journal v2 entry is missing its backend identity"))?,
|
.ok_or_else(|| Error::other("tier delete journal v2 entry is missing its backend identity"))?,
|
||||||
),
|
),
|
||||||
false,
|
false,
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown,
|
||||||
),
|
),
|
||||||
TIER_DELETE_JOURNAL_EXACT_VERSION => {
|
TIER_DELETE_JOURNAL_EXACT_VERSION => {
|
||||||
if self.version_id.is_empty() || self.version_id_exact != Some(true) {
|
if self.version_id.is_empty() || self.version_id_exact != Some(true) {
|
||||||
@@ -108,6 +127,22 @@ impl PersistedTierDeleteJournalEntry {
|
|||||||
.ok_or_else(|| Error::other("tier delete journal v3 entry is missing its backend identity"))?,
|
.ok_or_else(|| Error::other("tier delete journal v3 entry is missing its backend identity"))?,
|
||||||
),
|
),
|
||||||
true,
|
true,
|
||||||
|
rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
TIER_DELETE_JOURNAL_STATE_VERSION => {
|
||||||
|
let state = self
|
||||||
|
.version_state
|
||||||
|
.ok_or_else(|| Error::other("tier delete journal v4 entry is missing its version state"))?;
|
||||||
|
let exact = self.version_id_exact.unwrap_or(false);
|
||||||
|
validate_version_state(state, &self.version_id, exact)?;
|
||||||
|
(
|
||||||
|
Some(
|
||||||
|
self.backend_identity
|
||||||
|
.ok_or_else(|| Error::other("tier delete journal v4 entry is missing its backend identity"))?,
|
||||||
|
),
|
||||||
|
exact,
|
||||||
|
state,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
version => return Err(Error::other(format!("unsupported tier delete journal version {version}"))),
|
version => return Err(Error::other(format!("unsupported tier delete journal version {version}"))),
|
||||||
@@ -118,10 +153,30 @@ impl PersistedTierDeleteJournalEntry {
|
|||||||
tier_name: self.tier_name,
|
tier_name: self.tier_name,
|
||||||
backend_identity,
|
backend_identity,
|
||||||
version_id_exact,
|
version_id_exact,
|
||||||
|
version_state,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn validate_version_state(
|
||||||
|
state: rustfs_filemeta::TransitionVersionState,
|
||||||
|
version_id: &str,
|
||||||
|
version_id_exact: bool,
|
||||||
|
) -> Result<()> {
|
||||||
|
use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown};
|
||||||
|
|
||||||
|
let valid = match state {
|
||||||
|
Unknown => !version_id_exact,
|
||||||
|
KnownDisabled => version_id.is_empty() && !version_id_exact,
|
||||||
|
SuspendedNull => version_id == "null" && version_id_exact,
|
||||||
|
Exact => !version_id.is_empty() && version_id != "null" && version_id_exact,
|
||||||
|
};
|
||||||
|
if !valid {
|
||||||
|
return Err(Error::other("tier delete journal version state conflicts with its version id"));
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct TierDeleteJournalRecoveryStats {
|
pub struct TierDeleteJournalRecoveryStats {
|
||||||
pub scanned: usize,
|
pub scanned: usize,
|
||||||
@@ -159,7 +214,7 @@ pub(crate) fn decode_tier_delete_journal_entry(data: &[u8]) -> Result<Jentry> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn encode_tier_delete_journal_entry(je: &Jentry) -> Result<Vec<u8>> {
|
pub(crate) fn encode_tier_delete_journal_entry(je: &Jentry) -> Result<Vec<u8>> {
|
||||||
serde_json::to_vec(&PersistedTierDeleteJournalEntry::from_jentry(je))
|
serde_json::to_vec(&PersistedTierDeleteJournalEntry::from_jentry(je)?)
|
||||||
.map_err(|err| Error::other(format!("encode tier delete journal failed: {err}")))
|
.map_err(|err| Error::other(format!("encode tier delete journal failed: {err}")))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -209,18 +264,35 @@ where
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn process_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jentry) -> std::io::Result<()> {
|
pub async fn process_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jentry) -> std::io::Result<()> {
|
||||||
|
if je.version_state == rustfs_filemeta::TransitionVersionState::Unknown {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"tier delete journal remote version state is unknown",
|
||||||
|
));
|
||||||
|
}
|
||||||
let backend_identity = je
|
let backend_identity = je
|
||||||
.backend_identity
|
.backend_identity
|
||||||
.ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?;
|
.ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?;
|
||||||
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
if je.version_id_exact {
|
||||||
&je.obj_name,
|
delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
&je.version_id,
|
&je.obj_name,
|
||||||
&je.tier_name,
|
&je.version_id,
|
||||||
backend_identity,
|
&je.tier_name,
|
||||||
&api.tier_config_mgr(),
|
backend_identity,
|
||||||
je.version_id_exact,
|
&api.tier_config_mgr(),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
} else {
|
||||||
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
|
||||||
|
&je.obj_name,
|
||||||
|
&je.version_id,
|
||||||
|
&je.tier_name,
|
||||||
|
backend_identity,
|
||||||
|
&api.tier_config_mgr(),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
remove_tier_delete_journal_entry(api, je).await
|
remove_tier_delete_journal_entry(api, je).await
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -406,8 +478,9 @@ where
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
TIER_DELETE_JOURNAL_EXACT_VERSION, await_tier_delete_journal_recovery, decode_tier_delete_journal_entry,
|
TIER_DELETE_JOURNAL_EXACT_VERSION, TIER_DELETE_JOURNAL_STATE_VERSION, await_tier_delete_journal_recovery,
|
||||||
encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity, tier_delete_journal_object_name,
|
decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity,
|
||||||
|
tier_delete_journal_object_name,
|
||||||
};
|
};
|
||||||
use crate::bucket::lifecycle::tier_sweeper::Jentry;
|
use crate::bucket::lifecycle::tier_sweeper::Jentry;
|
||||||
use crate::error::Result;
|
use crate::error::Result;
|
||||||
@@ -420,7 +493,8 @@ mod tests {
|
|||||||
version_id: "remote-version".to_string(),
|
version_id: "remote-version".to_string(),
|
||||||
tier_name: "WARM".to_string(),
|
tier_name: "WARM".to_string(),
|
||||||
backend_identity: Some([7; 32]),
|
backend_identity: Some([7; 32]),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -436,6 +510,7 @@ mod tests {
|
|||||||
assert_eq!(decoded.tier_name, je.tier_name);
|
assert_eq!(decoded.tier_name, je.tier_name);
|
||||||
assert_eq!(decoded.backend_identity, je.backend_identity);
|
assert_eq!(decoded.backend_identity, je.backend_identity);
|
||||||
assert_eq!(decoded.version_id_exact, je.version_id_exact);
|
assert_eq!(decoded.version_id_exact, je.version_id_exact);
|
||||||
|
assert_eq!(decoded.version_state, je.version_state);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -450,7 +525,7 @@ mod tests {
|
|||||||
let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("exact journal JSON should decode");
|
let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("exact journal JSON should decode");
|
||||||
let decoded = decode_tier_delete_journal_entry(&encoded).expect("exact journal entry should decode");
|
let decoded = decode_tier_delete_journal_entry(&encoded).expect("exact journal entry should decode");
|
||||||
|
|
||||||
assert_eq!(persisted["version"], TIER_DELETE_JOURNAL_EXACT_VERSION);
|
assert_eq!(persisted["version"], TIER_DELETE_JOURNAL_STATE_VERSION);
|
||||||
assert_eq!(persisted["version_id_exact"], true);
|
assert_eq!(persisted["version_id_exact"], true);
|
||||||
assert!(decoded.version_id_exact);
|
assert!(decoded.version_id_exact);
|
||||||
assert_ne!(tier_delete_journal_object_name(&exact), tier_delete_journal_object_name(&normalized));
|
assert_ne!(tier_delete_journal_object_name(&exact), tier_delete_journal_object_name(&normalized));
|
||||||
@@ -513,6 +588,46 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn tier_delete_journal_rejects_conflicting_v4_version_states() {
|
||||||
|
let identity = vec![7_u8; 32];
|
||||||
|
let invalid = [
|
||||||
|
("known-disabled", "unexpected", false),
|
||||||
|
("suspended-null", "", true),
|
||||||
|
("suspended-null", "null", false),
|
||||||
|
("exact", "", true),
|
||||||
|
("exact", "null", true),
|
||||||
|
("exact", "version", false),
|
||||||
|
("unknown", "version", true),
|
||||||
|
];
|
||||||
|
|
||||||
|
for (state, version_id, exact) in invalid {
|
||||||
|
let persisted = serde_json::json!({
|
||||||
|
"version": TIER_DELETE_JOURNAL_STATE_VERSION,
|
||||||
|
"obj_name": "remote/object",
|
||||||
|
"version_id": version_id,
|
||||||
|
"tier_name": "WARM",
|
||||||
|
"backend_identity": identity,
|
||||||
|
"version_id_exact": exact.then_some(true),
|
||||||
|
"version_state": state,
|
||||||
|
});
|
||||||
|
let encoded = serde_json::to_vec(&persisted).expect("invalid journal fixture should encode");
|
||||||
|
decode_tier_delete_journal_entry(&encoded).expect_err("conflicting v4 version state must fail closed");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn legacy_journals_decode_with_unknown_version_state() {
|
||||||
|
let v1 = br#"{"version":1,"obj_name":"remote/object","version_id":"opaque","tier_name":"WARM"}"#;
|
||||||
|
let v2 = br#"{"version":2,"obj_name":"remote/object","version_id":"opaque","tier_name":"WARM","backend_identity":[7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7]}"#;
|
||||||
|
|
||||||
|
for payload in [v1.as_slice(), v2.as_slice()] {
|
||||||
|
let decoded = decode_tier_delete_journal_entry(payload).expect("legacy journal should decode");
|
||||||
|
assert_eq!(decoded.version_state, rustfs_filemeta::TransitionVersionState::Unknown);
|
||||||
|
assert!(!decoded.version_id_exact);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn tier_delete_journal_path_is_stable_and_sanitized() {
|
fn tier_delete_journal_path_is_stable_and_sanitized() {
|
||||||
let je = journal_entry();
|
let je = journal_entry();
|
||||||
@@ -530,6 +645,8 @@ mod tests {
|
|||||||
fn tier_delete_journal_paths_separate_legacy_and_backend_identities() {
|
fn tier_delete_journal_paths_separate_legacy_and_backend_identities() {
|
||||||
let mut legacy = journal_entry();
|
let mut legacy = journal_entry();
|
||||||
legacy.backend_identity = None;
|
legacy.backend_identity = None;
|
||||||
|
legacy.version_id_exact = false;
|
||||||
|
legacy.version_state = rustfs_filemeta::TransitionVersionState::Unknown;
|
||||||
let mut backend_a = journal_entry();
|
let mut backend_a = journal_entry();
|
||||||
backend_a.backend_identity = Some([1; 32]);
|
backend_a.backend_identity = Some([1; 32]);
|
||||||
let mut backend_b = journal_entry();
|
let mut backend_b = journal_entry();
|
||||||
@@ -575,6 +692,8 @@ mod tests {
|
|||||||
fn tier_delete_journal_without_transition_identity_stays_legacy() {
|
fn tier_delete_journal_without_transition_identity_stays_legacy() {
|
||||||
let mut je = journal_entry();
|
let mut je = journal_entry();
|
||||||
je.backend_identity = None;
|
je.backend_identity = None;
|
||||||
|
je.version_id_exact = false;
|
||||||
|
je.version_state = rustfs_filemeta::TransitionVersionState::Unknown;
|
||||||
|
|
||||||
let encoded = encode_tier_delete_journal_entry(&je).expect("legacy journal should remain encodable");
|
let encoded = encode_tier_delete_journal_entry(&je).expect("legacy journal should remain encodable");
|
||||||
let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("journal JSON should decode");
|
let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("journal JSON should decode");
|
||||||
|
|||||||
@@ -185,6 +185,7 @@ struct ObjSweeper {
|
|||||||
transition_status: String,
|
transition_status: String,
|
||||||
transition_tier: String,
|
transition_tier: String,
|
||||||
transition_version_id: String,
|
transition_version_id: String,
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
remote_object: String,
|
remote_object: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -231,7 +232,9 @@ impl ObjSweeper {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn should_remove_remote_object(&self) -> Option<Jentry> {
|
pub fn should_remove_remote_object(&self) -> Option<Jentry> {
|
||||||
if self.transition_status != lifecycle::TRANSITION_COMPLETE {
|
if self.transition_status != lifecycle::TRANSITION_COMPLETE
|
||||||
|
|| self.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
|
{
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -249,7 +252,11 @@ impl ObjSweeper {
|
|||||||
version_id: self.transition_version_id.clone(),
|
version_id: self.transition_version_id.clone(),
|
||||||
tier_name: self.transition_tier.clone(),
|
tier_name: self.transition_tier.clone(),
|
||||||
backend_identity: None,
|
backend_identity: None,
|
||||||
version_id_exact: false,
|
version_id_exact: matches!(
|
||||||
|
self.transition_version_state,
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact
|
||||||
|
),
|
||||||
|
version_state: self.transition_version_state,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
None
|
None
|
||||||
@@ -286,6 +293,7 @@ pub struct Jentry {
|
|||||||
pub(crate) tier_name: String,
|
pub(crate) tier_name: String,
|
||||||
pub(crate) backend_identity: Option<TierDestinationId>,
|
pub(crate) backend_identity: Option<TierDestinationId>,
|
||||||
pub(crate) version_id_exact: bool,
|
pub(crate) version_id_exact: bool,
|
||||||
|
pub(crate) version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ExpiryOp for Jentry {
|
impl ExpiryOp for Jentry {
|
||||||
@@ -330,7 +338,7 @@ async fn delete_object_from_remote_tier_raw_with_manager(
|
|||||||
let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name)
|
let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name)
|
||||||
.await
|
.await
|
||||||
.map_err(std::io::Error::other)?;
|
.map_err(std::io::Error::other)?;
|
||||||
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false).await
|
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false, true).await
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn delete_object_from_remote_tier_raw_with_lease(
|
async fn delete_object_from_remote_tier_raw_with_lease(
|
||||||
@@ -338,8 +346,11 @@ async fn delete_object_from_remote_tier_raw_with_lease(
|
|||||||
rv_id: &str,
|
rv_id: &str,
|
||||||
lease: &TierOperationLease,
|
lease: &TierOperationLease,
|
||||||
version_id_exact: bool,
|
version_id_exact: bool,
|
||||||
|
validate_remote_version_id: bool,
|
||||||
) -> Result<(), std::io::Error> {
|
) -> Result<(), std::io::Error> {
|
||||||
lease.validate_remote_version_id(rv_id)?;
|
if validate_remote_version_id {
|
||||||
|
lease.validate_remote_version_id(rv_id)?;
|
||||||
|
}
|
||||||
|
|
||||||
if remote_delete_breaker_is_open(Instant::now()).await {
|
if remote_delete_breaker_is_open(Instant::now()).await {
|
||||||
metrics::counter!(METRIC_DELETE_REMOTE_BREAKER_TOTAL).increment(1);
|
metrics::counter!(METRIC_DELETE_REMOTE_BREAKER_TOTAL).increment(1);
|
||||||
@@ -435,7 +446,53 @@ pub(crate) async fn delete_object_from_remote_tier_with_lease_idempotent(
|
|||||||
lease: &TierOperationLease,
|
lease: &TierOperationLease,
|
||||||
version_id_exact: bool,
|
version_id_exact: bool,
|
||||||
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
|
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
|
||||||
match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact).await {
|
delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, version_id_exact, true).await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn delete_confirmed_transition_candidate_exact_with_lease_idempotent(
|
||||||
|
obj_name: &str,
|
||||||
|
rv_id: &str,
|
||||||
|
lease: &TierOperationLease,
|
||||||
|
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
|
||||||
|
if rv_id.is_empty() {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidInput,
|
||||||
|
"confirmed versioned transition candidate requires a non-empty remote version",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
#[cfg(test)]
|
||||||
|
if obj_name == "remote/empty-guard-probe" {
|
||||||
|
CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, true, false).await
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
|
||||||
|
|
||||||
|
pub(crate) async fn delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
|
obj_name: &str,
|
||||||
|
rv_id: &str,
|
||||||
|
tier_name: &str,
|
||||||
|
backend_identity: TierDestinationId,
|
||||||
|
tier_config_mgr: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
|
||||||
|
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
|
||||||
|
let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, tier_name, backend_identity)
|
||||||
|
.await
|
||||||
|
.map_err(std::io::Error::other)?;
|
||||||
|
delete_confirmed_transition_candidate_exact_with_lease_idempotent(obj_name, rv_id, &lease).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_object_from_remote_tier_with_lease_idempotent_inner(
|
||||||
|
obj_name: &str,
|
||||||
|
rv_id: &str,
|
||||||
|
lease: &TierOperationLease,
|
||||||
|
version_id_exact: bool,
|
||||||
|
validate_remote_version_id: bool,
|
||||||
|
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
|
||||||
|
match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact, validate_remote_version_id)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(()) => Ok(RemoteTierDeleteOutcome::Deleted),
|
Ok(()) => Ok(RemoteTierDeleteOutcome::Deleted),
|
||||||
Err(err) if is_remote_tier_not_found_error(&err) => Ok(RemoteTierDeleteOutcome::AlreadyRemoved),
|
Err(err) if is_remote_tier_not_found_error(&err) => Ok(RemoteTierDeleteOutcome::AlreadyRemoved),
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -460,6 +517,7 @@ pub fn transitioned_delete_journal_entry(
|
|||||||
versioned: bool,
|
versioned: bool,
|
||||||
suspended: bool,
|
suspended: bool,
|
||||||
transitioned: &TransitionedObject,
|
transitioned: &TransitionedObject,
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
) -> Option<Jentry> {
|
) -> Option<Jentry> {
|
||||||
let sweeper = ObjSweeper {
|
let sweeper = ObjSweeper {
|
||||||
version_id,
|
version_id,
|
||||||
@@ -468,6 +526,7 @@ pub fn transitioned_delete_journal_entry(
|
|||||||
transition_status: transitioned.status.clone(),
|
transition_status: transitioned.status.clone(),
|
||||||
transition_tier: transitioned.tier.clone(),
|
transition_tier: transitioned.tier.clone(),
|
||||||
transition_version_id: transitioned.version_id.clone(),
|
transition_version_id: transitioned.version_id.clone(),
|
||||||
|
transition_version_state,
|
||||||
remote_object: transitioned.name.clone(),
|
remote_object: transitioned.name.clone(),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
@@ -475,8 +534,13 @@ pub fn transitioned_delete_journal_entry(
|
|||||||
sweeper.should_remove_remote_object()
|
sweeper.should_remove_remote_object()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject) -> Option<Jentry> {
|
pub fn transitioned_force_delete_journal_entry(
|
||||||
if transitioned.status != lifecycle::TRANSITION_COMPLETE {
|
transitioned: &TransitionedObject,
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
|
) -> Option<Jentry> {
|
||||||
|
if transitioned.status != lifecycle::TRANSITION_COMPLETE
|
||||||
|
|| transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
|
||||||
|
{
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -485,7 +549,11 @@ pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject
|
|||||||
version_id: transitioned.version_id.clone(),
|
version_id: transitioned.version_id.clone(),
|
||||||
tier_name: transitioned.tier.clone(),
|
tier_name: transitioned.tier.clone(),
|
||||||
backend_identity: None,
|
backend_identity: None,
|
||||||
version_id_exact: false,
|
version_id_exact: matches!(
|
||||||
|
transition_version_state,
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact
|
||||||
|
),
|
||||||
|
version_state: transition_version_state,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -494,11 +562,14 @@ mod test {
|
|||||||
use crate::client::signer_error::invalid_utf8_header_error;
|
use crate::client::signer_error::invalid_utf8_header_error;
|
||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, RemoteDeleteBreaker, RemoteTierDeleteOutcome,
|
CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES, ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED,
|
||||||
|
RemoteDeleteBreaker, RemoteTierDeleteOutcome, delete_confirmed_transition_candidate_exact_with_manager_and_identity,
|
||||||
delete_object_from_remote_tier_idempotent, delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
delete_object_from_remote_tier_idempotent, delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
||||||
is_remote_tier_not_found_error, is_signer_header_error, set_remote_tier_delete_test_hook,
|
is_remote_tier_not_found_error, is_signer_header_error, lifecycle, set_remote_tier_delete_test_hook,
|
||||||
should_record_remote_delete_failure,
|
should_record_remote_delete_failure, transitioned_delete_journal_entry, transitioned_force_delete_journal_entry,
|
||||||
};
|
};
|
||||||
|
use crate::storage_api_contracts::lifecycle::TransitionedObject;
|
||||||
|
use rustfs_filemeta::TransitionVersionState;
|
||||||
use std::io::{Error, ErrorKind};
|
use std::io::{Error, ErrorKind};
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
@@ -542,6 +613,43 @@ mod test {
|
|||||||
assert!(should_record_remote_delete_failure(&Error::other("NoSuchVersion")));
|
assert!(should_record_remote_delete_failure(&Error::other("NoSuchVersion")));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn transitioned_delete_journal_preserves_remote_version_state() {
|
||||||
|
let cases = [
|
||||||
|
(TransitionVersionState::Unknown, "legacy-version", None),
|
||||||
|
(TransitionVersionState::KnownDisabled, "", Some(false)),
|
||||||
|
(TransitionVersionState::SuspendedNull, "null", Some(true)),
|
||||||
|
(TransitionVersionState::Exact, "opaque-version", Some(true)),
|
||||||
|
];
|
||||||
|
|
||||||
|
for (state, version_id, expected_exact) in cases {
|
||||||
|
let transitioned = TransitionedObject {
|
||||||
|
name: "remote/object".to_string(),
|
||||||
|
version_id: version_id.to_string(),
|
||||||
|
tier: "WARM".to_string(),
|
||||||
|
status: lifecycle::TRANSITION_COMPLETE.to_string(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let regular = transitioned_delete_journal_entry(None, false, false, &transitioned, state);
|
||||||
|
let forced = transitioned_force_delete_journal_entry(&transitioned, state);
|
||||||
|
|
||||||
|
match expected_exact {
|
||||||
|
Some(expected_exact) => {
|
||||||
|
let regular = regular.expect("known version state should produce a regular delete journal entry");
|
||||||
|
assert_eq!(regular.version_state, state);
|
||||||
|
assert_eq!(regular.version_id_exact, expected_exact);
|
||||||
|
let forced = forced.expect("known version state should produce a forced delete journal entry");
|
||||||
|
assert_eq!(forced.version_state, state);
|
||||||
|
assert_eq!(forced.version_id_exact, expected_exact);
|
||||||
|
}
|
||||||
|
None => {
|
||||||
|
assert!(regular.is_none());
|
||||||
|
assert!(forced.is_none());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn idempotent_remote_delete_treats_hooked_nosuchversion_as_already_removed() {
|
async fn idempotent_remote_delete_treats_hooked_nosuchversion_as_already_removed() {
|
||||||
@@ -664,6 +772,55 @@ mod test {
|
|||||||
assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]);
|
assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn confirmed_transition_cleanup_deletes_exact_provider_token() {
|
||||||
|
CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.store(0, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
let manager = crate::services::tier::tier::TierConfigMgr::new();
|
||||||
|
let backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
|
||||||
|
let lease = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
|
||||||
|
.await
|
||||||
|
.expect("test tier lease should be available");
|
||||||
|
let identity = lease.backend_identity();
|
||||||
|
drop(lease);
|
||||||
|
backend.set_reject_non_empty_remote_versions(true);
|
||||||
|
|
||||||
|
let outcome = delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
|
"remote/object",
|
||||||
|
"provider-version-token",
|
||||||
|
"WARM",
|
||||||
|
identity,
|
||||||
|
&manager,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("confirmed upload compensation should delete the exact provider token");
|
||||||
|
|
||||||
|
assert_eq!(outcome, RemoteTierDeleteOutcome::Deleted);
|
||||||
|
assert_eq!(backend.exact_remove_count(), 1);
|
||||||
|
assert_eq!(
|
||||||
|
backend.remove_versions().await,
|
||||||
|
vec![("remote/object".to_string(), "provider-version-token".to_string())]
|
||||||
|
);
|
||||||
|
|
||||||
|
let err = delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
|
"remote/empty-guard-probe",
|
||||||
|
"",
|
||||||
|
"WARM",
|
||||||
|
identity,
|
||||||
|
&manager,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect_err("confirmed versioned cleanup must reject an empty token");
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
|
||||||
|
assert_eq!(backend.remove_count().await, 1);
|
||||||
|
assert_eq!(
|
||||||
|
CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.load(std::sync::atomic::Ordering::Relaxed),
|
||||||
|
0,
|
||||||
|
"empty remote versions must be rejected before exact cleanup dispatch"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn breaker_opens_at_threshold_and_recovers_after_window() {
|
fn breaker_opens_at_threshold_and_recovers_after_window() {
|
||||||
let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30));
|
let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30));
|
||||||
|
|||||||
@@ -22,7 +22,10 @@ use uuid::Uuid;
|
|||||||
|
|
||||||
use crate::bucket::lifecycle::config_boundary;
|
use crate::bucket::lifecycle::config_boundary;
|
||||||
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
||||||
use crate::bucket::lifecycle::tier_sweeper::delete_object_from_remote_tier_idempotent_with_manager_and_identity;
|
use crate::bucket::lifecycle::tier_sweeper::{
|
||||||
|
delete_confirmed_transition_candidate_exact_with_manager_and_identity,
|
||||||
|
delete_object_from_remote_tier_idempotent_with_manager_and_identity,
|
||||||
|
};
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
use crate::disk::RUSTFS_META_BUCKET;
|
||||||
use crate::error::{Error, Result as EcstoreResult};
|
use crate::error::{Error, Result as EcstoreResult};
|
||||||
use crate::object_api::ObjectOptions;
|
use crate::object_api::ObjectOptions;
|
||||||
@@ -708,6 +711,21 @@ async fn recover_unknown_upload_outcome(
|
|||||||
TransitionCandidateProbe::UnversionedPresent => {
|
TransitionCandidateProbe::UnversionedPresent => {
|
||||||
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
|
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await
|
||||||
}
|
}
|
||||||
|
TransitionCandidateProbe::VersionedPresent(version_id)
|
||||||
|
if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) =>
|
||||||
|
{
|
||||||
|
delete_confirmed_transition_candidate_exact_with_manager_and_identity(
|
||||||
|
&transaction.remote_object,
|
||||||
|
&version_id,
|
||||||
|
&transaction.tier_name,
|
||||||
|
transaction.backend_fingerprint,
|
||||||
|
&api.tier_config_mgr(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.map_err(Error::other)?;
|
||||||
|
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||||
|
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||||
|
}
|
||||||
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
||||||
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
|
cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -14,7 +14,7 @@
|
|||||||
|
|
||||||
use crate::cluster::rpc::http_auth::RPC_CONTENT_SHA256_HEADER;
|
use crate::cluster::rpc::http_auth::RPC_CONTENT_SHA256_HEADER;
|
||||||
use crate::cluster::rpc::{gen_tonic_signature_headers, normalize_tonic_rpc_audience};
|
use crate::cluster::rpc::{gen_tonic_signature_headers, normalize_tonic_rpc_audience};
|
||||||
use crate::disk::error::{DiskError, Error as DiskErrorType};
|
use crate::disk::error::{DiskError, Error as DiskErrorType, RpcStatusError};
|
||||||
use crate::runtime::sources as runtime_sources;
|
use crate::runtime::sources as runtime_sources;
|
||||||
use http::Uri;
|
use http::Uri;
|
||||||
use rustfs_protos::{
|
use rustfs_protos::{
|
||||||
@@ -107,10 +107,79 @@ pub async fn node_service_time_out_client_no_auth(
|
|||||||
node_service_time_out_client(addr, TonicInterceptor::NoOp(NoOpInterceptor)).await
|
node_service_time_out_client(addr, TonicInterceptor::NoOp(NoOpInterceptor)).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The typed `tonic::Status` an internode RPC failure was converted from, if
|
||||||
|
/// this error carries one.
|
||||||
|
pub(crate) fn embedded_tonic_status(io_err: &std::io::Error) -> Option<&tonic::Status> {
|
||||||
|
io_err.get_ref()?.downcast_ref::<RpcStatusError>().map(RpcStatusError::status)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Decide whether a gRPC status reports a peer we cannot currently reach,
|
||||||
|
/// rather than an application outcome from a live peer.
|
||||||
|
///
|
||||||
|
/// `Unavailable` is the one code that means "no service behind this channel":
|
||||||
|
/// the client transport raises it when the connection is broken, and the
|
||||||
|
/// server's own not-ready gates use it deliberately.
|
||||||
|
///
|
||||||
|
/// `Unknown` is the client transport's escape hatch for a cause it could not
|
||||||
|
/// map to a code — tower's "Service was not ready: <cause>", an h2 error with
|
||||||
|
/// no gRPC mapping. Our handlers never return it, so there its message is the
|
||||||
|
/// only evidence available and the anchored needles decide.
|
||||||
|
///
|
||||||
|
/// Every other code is an answer from a live peer and is never a transport
|
||||||
|
/// failure, whatever its message says. That distinction is the point of
|
||||||
|
/// classifying by code: a peer relaying its own downstream trouble as
|
||||||
|
/// `Internal("connection refused ...")`, or a handler interpolating a local
|
||||||
|
/// `io::Error` into `Status::internal`, answered us perfectly well. Marking it
|
||||||
|
/// offline over that text is the bug this classification replaces. Likewise a
|
||||||
|
/// `Cancelled` "Timeout expired" from the per-RPC channel deadline means the
|
||||||
|
/// peer is slow, not gone; gating it would turn load into a partition.
|
||||||
|
pub(crate) fn is_network_like_status(status: &tonic::Status) -> bool {
|
||||||
|
match status.code() {
|
||||||
|
tonic::Code::Unavailable => true,
|
||||||
|
tonic::Code::Unknown => message_has_network_needle(&status.to_string()),
|
||||||
|
_ => false,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Substring fallback for failures that only exist as text: dial errors
|
||||||
|
/// wrapped by `get_client`, remote `error_info` payloads, and statuses
|
||||||
|
/// flattened through `format!`. Needles must stay anchored to transport
|
||||||
|
/// context — a bare word like "unavailable" also matches application text
|
||||||
|
/// (e.g. a bucket named "unavailable-logs") and would take a healthy peer
|
||||||
|
/// offline.
|
||||||
|
pub(crate) fn message_has_network_needle(message: &str) -> bool {
|
||||||
|
let message = message.to_ascii_lowercase();
|
||||||
|
[
|
||||||
|
"temporarily offline",
|
||||||
|
"transport error",
|
||||||
|
// tonic >= 0.14 renders Code::Unavailable as
|
||||||
|
// `code: 'The service is currently unavailable'`.
|
||||||
|
"code: 'the service is currently unavailable'",
|
||||||
|
// RUSTFS_COMPAT_TODO(tonic-013-status-render): releases up to 1.0.0-alpha.38 shipped tonic 0.13, which rendered the same status as `status: Unavailable`, and peers relay that text in error_info. Remove after the minimum supported RustFS peer version ships tonic >= 0.14.
|
||||||
|
"status: unavailable",
|
||||||
|
"error trying to connect",
|
||||||
|
"connection refused",
|
||||||
|
"connection reset",
|
||||||
|
"broken pipe",
|
||||||
|
"not connected",
|
||||||
|
"unexpected eof",
|
||||||
|
"timed out",
|
||||||
|
"deadline has elapsed",
|
||||||
|
"connection closed",
|
||||||
|
"connection aborted",
|
||||||
|
"tcp connect error",
|
||||||
|
]
|
||||||
|
.iter()
|
||||||
|
.any(|needle| message.contains(needle))
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn is_network_like_disk_error(err: &DiskErrorType) -> bool {
|
pub(crate) fn is_network_like_disk_error(err: &DiskErrorType) -> bool {
|
||||||
match err {
|
match err {
|
||||||
DiskError::Timeout => true,
|
DiskError::Timeout => true,
|
||||||
DiskError::Io(io_err) => {
|
DiskError::Io(io_err) => {
|
||||||
|
if let Some(status) = embedded_tonic_status(io_err) {
|
||||||
|
return is_network_like_status(status);
|
||||||
|
}
|
||||||
if matches!(
|
if matches!(
|
||||||
io_err.kind(),
|
io_err.kind(),
|
||||||
ErrorKind::TimedOut
|
ErrorKind::TimedOut
|
||||||
@@ -124,24 +193,7 @@ pub(crate) fn is_network_like_disk_error(err: &DiskErrorType) -> bool {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
let message = io_err.to_string().to_ascii_lowercase();
|
message_has_network_needle(&io_err.to_string())
|
||||||
[
|
|
||||||
"transport error",
|
|
||||||
"unavailable",
|
|
||||||
"error trying to connect",
|
|
||||||
"connection refused",
|
|
||||||
"connection reset",
|
|
||||||
"broken pipe",
|
|
||||||
"not connected",
|
|
||||||
"unexpected eof",
|
|
||||||
"timed out",
|
|
||||||
"deadline has elapsed",
|
|
||||||
"connection closed",
|
|
||||||
"connection aborted",
|
|
||||||
"tcp connect error",
|
|
||||||
]
|
|
||||||
.iter()
|
|
||||||
.any(|needle| message.contains(needle))
|
|
||||||
}
|
}
|
||||||
_ => false,
|
_ => false,
|
||||||
}
|
}
|
||||||
@@ -269,6 +321,70 @@ mod tests {
|
|||||||
let _ = provider.shutdown();
|
let _ = provider.shutdown();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn network_like_disk_error_uses_typed_status_code() {
|
||||||
|
// Transport-level Unavailable statuses justify retry/eviction.
|
||||||
|
assert!(is_network_like_disk_error(&DiskError::from(tonic::Status::unavailable(
|
||||||
|
"storage layer is not initialized"
|
||||||
|
))));
|
||||||
|
// Application statuses from a live peer must not look network-like,
|
||||||
|
// even when their message contains transport-sounding words.
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::from(tonic::Status::internal(
|
||||||
|
"failed to heal bucket \"unavailable-logs\""
|
||||||
|
))));
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::from(tonic::Status::unauthenticated(
|
||||||
|
"No valid auth token"
|
||||||
|
))));
|
||||||
|
// A slow peer that blew the per-RPC deadline is still answering.
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::from(tonic::Status::cancelled("Timeout expired"))));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn embedded_tonic_status_is_recovered_across_error_conversions() {
|
||||||
|
// DiskError and StorageError share one wrapper, so a status keeps its
|
||||||
|
// typed classification whichever error it was converted into first.
|
||||||
|
let from_storage: DiskErrorType = crate::error::Error::from(tonic::Status::unavailable("peer gone")).into();
|
||||||
|
let DiskError::Io(io_err) = &from_storage else {
|
||||||
|
panic!("status-derived disk error should stay an Io error");
|
||||||
|
};
|
||||||
|
assert_eq!(embedded_tonic_status(io_err).map(|status| status.code()), Some(tonic::Code::Unavailable));
|
||||||
|
|
||||||
|
let from_disk = crate::error::Error::from(DiskError::from(tonic::Status::unavailable("peer gone")));
|
||||||
|
let crate::error::Error::Io(io_err) = &from_disk else {
|
||||||
|
panic!("status-derived storage error should stay an Io error");
|
||||||
|
};
|
||||||
|
assert_eq!(embedded_tonic_status(io_err).map(|status| status.code()), Some(tonic::Code::Unavailable));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn network_like_disk_error_ignores_transport_words_in_application_statuses() {
|
||||||
|
// Same contract as the peer client: a status the peer answered with
|
||||||
|
// is not a transport failure, so it must not drive a reconnect even
|
||||||
|
// when its message describes one.
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::from(tonic::Status::internal(
|
||||||
|
"connection refused while dialing downstream backend"
|
||||||
|
))));
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::from(tonic::Status::unauthenticated(
|
||||||
|
"connection reset while validating token"
|
||||||
|
))));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn network_like_disk_error_requires_anchored_unavailable_needle() {
|
||||||
|
// Regression: a bare "unavailable" needle used to match application
|
||||||
|
// text such as a bucket name.
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::other("bucket \"unavailable-logs\" not found")));
|
||||||
|
// Anchored renderings of a flattened Unavailable status still match.
|
||||||
|
assert!(is_network_like_disk_error(&DiskError::other(
|
||||||
|
"code: 'The service is currently unavailable', message: \"peer gone\""
|
||||||
|
)));
|
||||||
|
assert!(is_network_like_disk_error(&DiskError::other(
|
||||||
|
"status: Unavailable, message: \"peer gone\""
|
||||||
|
)));
|
||||||
|
assert!(is_network_like_disk_error(&DiskError::other("connection refused")));
|
||||||
|
assert!(!is_network_like_disk_error(&DiskError::FileNotFound));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_signature_interceptor_keeps_auth_headers() {
|
fn test_signature_interceptor_keeps_auth_headers() {
|
||||||
ensure_test_rpc_secret();
|
ensure_test_rpc_secret();
|
||||||
|
|||||||
@@ -13,8 +13,8 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use crate::cluster::rpc::client::{
|
use crate::cluster::rpc::client::{
|
||||||
TonicInterceptor, gen_tonic_signature_interceptor, heal_control_time_out_client, node_service_time_out_client,
|
TonicInterceptor, embedded_tonic_status, gen_tonic_signature_interceptor, heal_control_time_out_client,
|
||||||
tier_mutation_control_time_out_client,
|
is_network_like_status, message_has_network_needle, node_service_time_out_client, tier_mutation_control_time_out_client,
|
||||||
};
|
};
|
||||||
use crate::cluster::rpc::{set_tonic_canonical_body_digest, verify_tonic_rpc_response_proof};
|
use crate::cluster::rpc::{set_tonic_canonical_body_digest, verify_tonic_rpc_response_proof};
|
||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
@@ -471,26 +471,21 @@ impl PeerRestClient {
|
|||||||
self.offline.store(false, Ordering::Release);
|
self.offline.store(false, Ordering::Release);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Whether this failure means the peer is unreachable, so it should be
|
||||||
|
/// gated offline and its connection evicted.
|
||||||
|
///
|
||||||
|
/// RPC failures are classified by their typed gRPC code first
|
||||||
|
/// (`is_network_like_status`); an application error from a live peer must
|
||||||
|
/// never take it offline no matter what its message says. The substring
|
||||||
|
/// fallback only covers failures that exist purely as text, such as the
|
||||||
|
/// dial errors `get_client` wraps.
|
||||||
fn is_network_like_error(err: &Error) -> bool {
|
fn is_network_like_error(err: &Error) -> bool {
|
||||||
let message = err.to_string().to_ascii_lowercase();
|
if let Error::Io(io_err) = err
|
||||||
[
|
&& let Some(status) = embedded_tonic_status(io_err)
|
||||||
"temporarily offline",
|
{
|
||||||
"transport error",
|
return is_network_like_status(status);
|
||||||
"unavailable",
|
}
|
||||||
"error trying to connect",
|
message_has_network_needle(&err.to_string())
|
||||||
"connection refused",
|
|
||||||
"connection reset",
|
|
||||||
"broken pipe",
|
|
||||||
"not connected",
|
|
||||||
"unexpected eof",
|
|
||||||
"timed out",
|
|
||||||
"deadline has elapsed",
|
|
||||||
"connection closed",
|
|
||||||
"connection aborted",
|
|
||||||
"tcp connect error",
|
|
||||||
]
|
|
||||||
.iter()
|
|
||||||
.any(|needle| message.contains(needle))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn mark_offline_and_spawn_recovery(&self) {
|
fn mark_offline_and_spawn_recovery(&self) {
|
||||||
@@ -1702,30 +1697,17 @@ fn tier_config_reload_connection_outcome(err: Error) -> TierConfigReloadOutcome
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn is_tier_config_reload_connection_failure(err: &Error) -> bool {
|
fn is_tier_config_reload_connection_failure(err: &Error) -> bool {
|
||||||
let message = err.to_string().to_ascii_lowercase();
|
let message = err.to_string();
|
||||||
|
// A bare "unavailable" is only trusted inside the local dial-failure
|
||||||
|
// wrapper from `get_client`, never in application text.
|
||||||
if message
|
if message
|
||||||
|
.to_ascii_lowercase()
|
||||||
.split_once("can not get client, err:")
|
.split_once("can not get client, err:")
|
||||||
.is_some_and(|(_, local_error)| local_error.contains("unavailable"))
|
.is_some_and(|(_, local_error)| local_error.contains("unavailable"))
|
||||||
{
|
{
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
[
|
message_has_network_needle(&message)
|
||||||
"temporarily offline",
|
|
||||||
"transport error",
|
|
||||||
"error trying to connect",
|
|
||||||
"connection refused",
|
|
||||||
"connection reset",
|
|
||||||
"connection closed",
|
|
||||||
"connection aborted",
|
|
||||||
"broken pipe",
|
|
||||||
"not connected",
|
|
||||||
"unexpected eof",
|
|
||||||
"timed out",
|
|
||||||
"deadline has elapsed",
|
|
||||||
"tcp connect error",
|
|
||||||
]
|
|
||||||
.iter()
|
|
||||||
.any(|needle| message.contains(needle))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn tier_config_reload_remote_failure(error_info: Option<String>) -> TierConfigReloadOutcome {
|
fn tier_config_reload_remote_failure(error_info: Option<String>) -> TierConfigReloadOutcome {
|
||||||
@@ -2151,6 +2133,126 @@ mod tests {
|
|||||||
assert!(!PeerRestClient::is_network_like_error(&Error::NotImplemented));
|
assert!(!PeerRestClient::is_network_like_error(&Error::NotImplemented));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn peer_rest_client_network_classifier_uses_typed_status_code() {
|
||||||
|
// The one code that means "nothing is answering on this channel".
|
||||||
|
assert!(PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unavailable(
|
||||||
|
"storage layer is not initialized"
|
||||||
|
))));
|
||||||
|
// Application statuses from a live peer must not mark it offline,
|
||||||
|
// even when their message contains transport-sounding words.
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::internal(
|
||||||
|
"failed to reload metadata for bucket \"unavailable-logs\""
|
||||||
|
))));
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unauthenticated(
|
||||||
|
"No valid auth token"
|
||||||
|
))));
|
||||||
|
// A request-budget expiry answered by a live peer is an application
|
||||||
|
// outcome, not a transport failure.
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::deadline_exceeded(
|
||||||
|
"heal control request expired"
|
||||||
|
))));
|
||||||
|
// Unknown is the transport's escape hatch for a cause it could not
|
||||||
|
// map, and our handlers never return it, so there the text decides.
|
||||||
|
assert!(PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unknown(
|
||||||
|
"Service was not ready: transport error"
|
||||||
|
))));
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unknown(
|
||||||
|
"peer response unknown"
|
||||||
|
))));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn peer_rest_client_network_classifier_ignores_transport_words_in_application_statuses() {
|
||||||
|
// The reason classification reads the code rather than the text: a
|
||||||
|
// peer that answers is reachable, even when what it says describes a
|
||||||
|
// connection failure of its own. A handler interpolating a local
|
||||||
|
// io::Error into Status::internal, or relaying trouble with its own
|
||||||
|
// downstream, must not cost us the channel to a healthy peer.
|
||||||
|
for status in [
|
||||||
|
tonic::Status::internal("connection refused while dialing downstream backend"),
|
||||||
|
tonic::Status::internal("write failed: broken pipe"),
|
||||||
|
tonic::Status::unauthenticated("connection reset while validating token"),
|
||||||
|
tonic::Status::failed_precondition("scanner lease timed out"),
|
||||||
|
tonic::Status::deadline_exceeded("heal control request timed out"),
|
||||||
|
] {
|
||||||
|
let rendered = status.to_string();
|
||||||
|
assert!(
|
||||||
|
!PeerRestClient::is_network_like_error(&Error::from(status)),
|
||||||
|
"an answered application status must not mark the peer offline: {rendered}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn peer_rest_client_network_classifier_keeps_slow_peers_online() {
|
||||||
|
// The per-RPC channel deadline (RUSTFS_INTERNODE_RPC_TIMEOUT, 30s)
|
||||||
|
// surfaces as Cancelled "Timeout expired" carrying the transport
|
||||||
|
// cause as its source. A peer that is merely slow must stay online:
|
||||||
|
// gating it would spend a full recovery cycle fast-failing every RPC
|
||||||
|
// to a host that is still answering, turning load into a partition.
|
||||||
|
let timeout_status = tonic::Status::cancelled("Timeout expired");
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(timeout_status)));
|
||||||
|
|
||||||
|
let sourced = tonic::Status::from_error(Box::new(std::io::Error::other("Timeout expired")));
|
||||||
|
assert!(
|
||||||
|
std::error::Error::source(&sourced).is_some(),
|
||||||
|
"the transport builds this status through Status::from_error, which attaches the cause"
|
||||||
|
);
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::from(sourced)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn rpc_status_errors_keep_their_rendering_and_hide_peer_metadata() {
|
||||||
|
let err = Error::from(tonic::Status::unavailable("peer gone"));
|
||||||
|
assert_eq!(
|
||||||
|
err.to_string(),
|
||||||
|
"Io error: code: 'The service is currently unavailable', message: \"peer gone\""
|
||||||
|
);
|
||||||
|
|
||||||
|
// tonic's own Debug prints the MetadataMap, i.e. every response header
|
||||||
|
// the peer sent; those must not reach a log through this error.
|
||||||
|
let mut status = tonic::Status::unavailable("peer gone");
|
||||||
|
status
|
||||||
|
.metadata_mut()
|
||||||
|
.insert("authorization", "Bearer secret".parse().expect("valid header value"));
|
||||||
|
let rendered = format!("{:?}", Error::from(status));
|
||||||
|
assert!(!rendered.contains("Bearer secret"), "{rendered}");
|
||||||
|
assert!(!rendered.contains("MetadataMap"), "{rendered}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn peer_rest_client_network_classifier_ignores_application_text_containing_unavailable() {
|
||||||
|
// Regression: a bare "unavailable" needle used to match application
|
||||||
|
// strings like these and take a healthy peer offline.
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::other(
|
||||||
|
"peer replication statistics provider is unavailable"
|
||||||
|
)));
|
||||||
|
assert!(!PeerRestClient::is_network_like_error(&Error::other(
|
||||||
|
"bucket \"unavailable-logs\" not found"
|
||||||
|
)));
|
||||||
|
// Anchored renderings of a flattened Unavailable status still match:
|
||||||
|
// tonic >= 0.14 form ...
|
||||||
|
assert!(PeerRestClient::is_network_like_error(&Error::other(
|
||||||
|
"peer tier mutation commit RPC failed: code: 'The service is currently unavailable', message: \"peer gone\""
|
||||||
|
)));
|
||||||
|
// ... which is only anchored as long as tonic renders Unavailable this
|
||||||
|
// way. A tonic bump that reworded it leaves the typed path correct but
|
||||||
|
// this needle stale, so pin the coupling rather than discover it in a
|
||||||
|
// partition.
|
||||||
|
assert!(
|
||||||
|
tonic::Status::unavailable("peer gone")
|
||||||
|
.to_string()
|
||||||
|
.to_ascii_lowercase()
|
||||||
|
.contains("code: 'the service is currently unavailable'"),
|
||||||
|
"tonic reworded Code::Unavailable; update the anchored needle"
|
||||||
|
);
|
||||||
|
// ... and the tonic <= 0.13 form peers may relay in error_info.
|
||||||
|
assert!(PeerRestClient::is_network_like_error(&Error::other(
|
||||||
|
"peer tier mutation commit RPC failed: status: Unavailable, message: \"peer gone\""
|
||||||
|
)));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn tier_config_reload_outcome_keeps_tonic_and_remote_errors_typed() {
|
fn tier_config_reload_outcome_keeps_tonic_and_remote_errors_typed() {
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
@@ -2193,6 +2295,14 @@ mod tests {
|
|||||||
tier_config_reload_connection_outcome(Error::other("can not get client, err: connection unavailable")),
|
tier_config_reload_connection_outcome(Error::other("can not get client, err: connection unavailable")),
|
||||||
TierConfigReloadOutcome::TransientReconnect(_)
|
TierConfigReloadOutcome::TransientReconnect(_)
|
||||||
));
|
));
|
||||||
|
// The bare word is trusted only to the right of the dial-failure
|
||||||
|
// prefix, not anywhere in the message.
|
||||||
|
assert!(matches!(
|
||||||
|
tier_config_reload_connection_outcome(Error::other(
|
||||||
|
"bucket unavailable-logs rejected it, then: can not get client, err: some other reason"
|
||||||
|
)),
|
||||||
|
TierConfigReloadOutcome::Terminal(_)
|
||||||
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -2495,6 +2605,39 @@ mod tests {
|
|||||||
assert!(!client.offline.load(Ordering::Acquire));
|
assert!(!client.offline.load(Ordering::Acquire));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn peer_rest_client_finalize_result_keeps_online_for_app_errors_mentioning_unavailable() {
|
||||||
|
// Regression: application error text containing "unavailable" (a
|
||||||
|
// remote error_info payload, or a bucket named "unavailable-logs" in
|
||||||
|
// a typed application status) must not take a healthy peer offline.
|
||||||
|
let client = test_peer_client();
|
||||||
|
let err = client
|
||||||
|
.finalize_result::<()>(Err(Error::other("peer replication statistics provider is unavailable")))
|
||||||
|
.await
|
||||||
|
.expect_err("application error should still be returned");
|
||||||
|
assert!(err.to_string().contains("provider is unavailable"));
|
||||||
|
assert!(!client.offline.load(Ordering::Acquire));
|
||||||
|
|
||||||
|
let err = client
|
||||||
|
.finalize_result::<()>(Err(Error::from(tonic::Status::internal(
|
||||||
|
"failed to reload metadata for bucket \"unavailable-logs\"",
|
||||||
|
))))
|
||||||
|
.await
|
||||||
|
.expect_err("application status should still be returned");
|
||||||
|
assert!(err.to_string().contains("unavailable-logs"));
|
||||||
|
assert!(!client.offline.load(Ordering::Acquire));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn peer_rest_client_finalize_result_marks_offline_for_typed_unavailable_status() {
|
||||||
|
let client = test_peer_client();
|
||||||
|
client
|
||||||
|
.finalize_result::<()>(Err(Error::from(tonic::Status::unavailable("storage layer is not initialized"))))
|
||||||
|
.await
|
||||||
|
.expect_err("network error should still be returned");
|
||||||
|
assert!(client.offline.load(Ordering::Acquire));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test(flavor = "current_thread")]
|
#[tokio::test(flavor = "current_thread")]
|
||||||
async fn peer_rest_recovery_probe_logs_keep_request_id_span_context() {
|
async fn peer_rest_recovery_probe_logs_keep_request_id_span_context() {
|
||||||
let logs = CapturedLogs::default();
|
let logs = CapturedLogs::default();
|
||||||
|
|||||||
@@ -2543,6 +2543,7 @@ mod tests {
|
|||||||
data_dir: None,
|
data_dir: None,
|
||||||
delete_marker: false,
|
delete_marker: false,
|
||||||
transitioned_object: Default::default(),
|
transitioned_object: Default::default(),
|
||||||
|
transition_version_state: Default::default(),
|
||||||
restore_ongoing: false,
|
restore_ongoing: false,
|
||||||
restore_expires: None,
|
restore_expires: None,
|
||||||
user_tags: Arc::new(String::new()),
|
user_tags: Arc::new(String::new()),
|
||||||
|
|||||||
@@ -13,9 +13,9 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use crate::disk::{
|
use crate::disk::{
|
||||||
CheckPartsResp, DeleteOptions, DiskAPI, DiskError, DiskInfo, DiskInfoOptions, DiskLocation, Endpoint, Error,
|
CheckPartsResp, DataDirDeleteStatus, DeleteOptions, DiskAPI, DiskError, DiskInfo, DiskInfoOptions, DiskLocation, Endpoint,
|
||||||
FileInfoVersions, MmapCopyStageMetrics, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, Result,
|
Error, FileInfoVersions, MmapCopyStageMetrics, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, Result,
|
||||||
UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
|
SnapshotLeaseToken, UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
|
||||||
health_state::{
|
health_state::{
|
||||||
RuntimeDriveHealthState, classify_drive_recovery, get_drive_returning_probe_interval,
|
RuntimeDriveHealthState, classify_drive_recovery, get_drive_returning_probe_interval,
|
||||||
get_drive_returning_success_threshold, get_drive_suspect_failure_threshold, record_drive_offline_duration,
|
get_drive_returning_success_threshold, get_drive_suspect_failure_threshold, record_drive_offline_duration,
|
||||||
@@ -1349,6 +1349,30 @@ impl DiskAPI for LocalDiskWrapper {
|
|||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result<SnapshotLeaseToken> {
|
||||||
|
self.track_disk_health(
|
||||||
|
|| async { self.disk.acquire_snapshot_lease(volume, path).await },
|
||||||
|
get_max_timeout_duration(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> {
|
||||||
|
self.track_disk_health(
|
||||||
|
|| async { self.disk.release_snapshot_lease(volume, path, token).await },
|
||||||
|
get_max_timeout_duration(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result<DataDirDeleteStatus> {
|
||||||
|
self.track_disk_health(
|
||||||
|
|| async { self.disk.delete_data_dir(volume, path, opts).await },
|
||||||
|
get_max_timeout_duration(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
|
async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
|
||||||
self.track_disk_health(
|
self.track_disk_health(
|
||||||
|| async { self.disk.write_metadata(org_volume, volume, path, fi).await },
|
|| async { self.disk.write_metadata(org_volume, volume, path, fi).await },
|
||||||
|
|||||||
@@ -337,9 +337,54 @@ impl From<DiskError> for std::io::Error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The single in-band representation of a failed internode RPC: it carries the
|
||||||
|
/// typed `tonic::Status` so failure classifiers can read the gRPC code instead
|
||||||
|
/// of substring-matching the rendered message (see `is_network_like_status`).
|
||||||
|
/// Both `DiskError` and `StorageError` wrap statuses in this type, so one
|
||||||
|
/// downcast recovers the status regardless of which error the status was
|
||||||
|
/// converted into first.
|
||||||
|
pub(crate) struct RpcStatusError(tonic::Status);
|
||||||
|
|
||||||
|
impl RpcStatusError {
|
||||||
|
pub(crate) fn status(&self) -> &tonic::Status {
|
||||||
|
&self.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<tonic::Status> for RpcStatusError {
|
||||||
|
fn from(status: tonic::Status) -> Self {
|
||||||
|
Self(status)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl std::fmt::Display for RpcStatusError {
|
||||||
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
|
std::fmt::Display::fmt(&self.0, f)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// `tonic::Status`'s own `Debug` prints its `MetadataMap`, i.e. every response
|
||||||
|
/// header and trailer the peer sent. Those are remote-controlled and can carry
|
||||||
|
/// credentials injected by a proxy in front of the peer, so keep them out of
|
||||||
|
/// anything that reaches a log.
|
||||||
|
impl std::fmt::Debug for RpcStatusError {
|
||||||
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
|
f.debug_struct("RpcStatusError")
|
||||||
|
.field("code", &self.0.code())
|
||||||
|
.field("message", &self.0.message())
|
||||||
|
.finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl StdError for RpcStatusError {
|
||||||
|
fn source(&self) -> Option<&(dyn StdError + 'static)> {
|
||||||
|
Some(&self.0)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl From<tonic::Status> for DiskError {
|
impl From<tonic::Status> for DiskError {
|
||||||
fn from(e: tonic::Status) -> Self {
|
fn from(e: tonic::Status) -> Self {
|
||||||
DiskError::other(e.message().to_string())
|
DiskError::Io(io::Error::other(RpcStatusError(e)))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -18,11 +18,12 @@ use crate::data_usage::local_snapshot::ensure_data_usage_layout;
|
|||||||
use crate::disk::disk_store::{get_drive_walkdir_stall_timeout, get_object_disk_read_timeout};
|
use crate::disk::disk_store::{get_drive_walkdir_stall_timeout, get_object_disk_read_timeout};
|
||||||
use crate::disk::{
|
use crate::disk::{
|
||||||
BUCKET_META_PREFIX, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN,
|
BUCKET_META_PREFIX, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN,
|
||||||
CHECK_PART_VOLUME_NOT_FOUND, CheckPartsResp, DeleteOptions, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskMetrics,
|
CHECK_PART_VOLUME_NOT_FOUND, CheckPartsResp, DataDirDeleteStatus, DeleteOptions, DiskAPI, DiskInfo, DiskInfoOptions,
|
||||||
FileInfoVersions, FileReader, FileWriter, MmapCopyStageMetrics, OldCurrentSize, PART_TRANSACTION_NEW_META,
|
DiskLocation, DiskMetrics, FileInfoVersions, FileReader, FileWriter, MmapCopyStageMetrics, OldCurrentSize,
|
||||||
PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction, RUSTFS_META_BUCKET, RUSTFS_META_TMP_BUCKET,
|
PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction, RUSTFS_META_BUCKET,
|
||||||
RUSTFS_META_TMP_DELETED_BUCKET, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, STORAGE_FORMAT_FILE,
|
RUSTFS_META_TMP_BUCKET, RUSTFS_META_TMP_DELETED_BUCKET, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp,
|
||||||
STORAGE_FORMAT_FILE_BACKUP, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, conv_part_err_to_int,
|
STORAGE_FORMAT_FILE, STORAGE_FORMAT_FILE_BACKUP, SnapshotLeaseToken, UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
|
||||||
|
conv_part_err_to_int,
|
||||||
endpoint::Endpoint,
|
endpoint::Endpoint,
|
||||||
error::{DiskError, Error, FileAccessDeniedWithContext, Result},
|
error::{DiskError, Error, FileAccessDeniedWithContext, Result},
|
||||||
error_conv::{to_access_error, to_file_error, to_unformatted_disk_error, to_volume_error},
|
error_conv::{to_access_error, to_file_error, to_unformatted_disk_error, to_volume_error},
|
||||||
@@ -66,7 +67,7 @@ use tokio::fs::{self, File};
|
|||||||
#[cfg(not(unix))]
|
#[cfg(not(unix))]
|
||||||
use tokio::io::AsyncReadExt;
|
use tokio::io::AsyncReadExt;
|
||||||
use tokio::io::{AsyncRead, AsyncSeekExt, AsyncWrite, AsyncWriteExt, ErrorKind, ReadBuf};
|
use tokio::io::{AsyncRead, AsyncSeekExt, AsyncWrite, AsyncWriteExt, ErrorKind, ReadBuf};
|
||||||
use tokio::sync::{Notify, RwLock, Semaphore};
|
use tokio::sync::{Mutex, Notify, RwLock, Semaphore};
|
||||||
use tokio::time::{Instant, Sleep, interval_at, timeout};
|
use tokio::time::{Instant, Sleep, interval_at, timeout};
|
||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, warn};
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
@@ -3856,6 +3857,25 @@ pub struct LocalDisk {
|
|||||||
exit_signal: Option<tokio::sync::broadcast::Sender<()>>,
|
exit_signal: Option<tokio::sync::broadcast::Sender<()>>,
|
||||||
io_backend: Arc<dyn LocalIoBackend>,
|
io_backend: Arc<dyn LocalIoBackend>,
|
||||||
file_sync_permits: Arc<Semaphore>,
|
file_sync_permits: Arc<Semaphore>,
|
||||||
|
snapshot_leases: Arc<Mutex<SnapshotLeaseRegistry>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
|
||||||
|
struct SnapshotLeaseKey {
|
||||||
|
volume: String,
|
||||||
|
path: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
struct SnapshotLeaseEntry {
|
||||||
|
tokens: HashSet<SnapshotLeaseToken>,
|
||||||
|
pending_delete: Option<DeleteOptions>,
|
||||||
|
deleting: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
struct SnapshotLeaseRegistry {
|
||||||
|
entries: HashMap<SnapshotLeaseKey, SnapshotLeaseEntry>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Drop for LocalDisk {
|
impl Drop for LocalDisk {
|
||||||
@@ -4092,6 +4112,7 @@ impl LocalDisk {
|
|||||||
exit_signal: None,
|
exit_signal: None,
|
||||||
io_backend: build_local_io_backend(root.clone()),
|
io_backend: build_local_io_backend(root.clone()),
|
||||||
file_sync_permits: os::disk_file_sync_limiter(&root),
|
file_sync_permits: os::disk_file_sync_limiter(&root),
|
||||||
|
snapshot_leases: Arc::new(Mutex::new(SnapshotLeaseRegistry::default())),
|
||||||
};
|
};
|
||||||
let (info, _root) = get_disk_info(root.clone()).await.inspect_err(|err| {
|
let (info, _root) = get_disk_info(root.clone()).await.inspect_err(|err| {
|
||||||
log_startup_disk_error("get_disk_info", &root, err);
|
log_startup_disk_error("get_disk_info", &root, err);
|
||||||
@@ -4576,6 +4597,23 @@ impl LocalDisk {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn delete_unleased(&self, volume: &str, path: &str, opt: &DeleteOptions) -> Result<()> {
|
||||||
|
let volume_dir = self.get_bucket_path(volume)?;
|
||||||
|
if !skip_access_checks(volume)
|
||||||
|
&& let Err(e) = access(&volume_dir).await
|
||||||
|
{
|
||||||
|
return Err(to_access_error(e, DiskError::VolumeAccessDenied).into());
|
||||||
|
}
|
||||||
|
|
||||||
|
let file_path = self.get_object_path(volume, path)?;
|
||||||
|
check_path_length(file_path.to_string_lossy().as_ref())?;
|
||||||
|
self.delete_file(&volume_dir, &file_path, opt.recursive, opt.immediate)
|
||||||
|
.await?;
|
||||||
|
// A deleted shard must not remain readable through the io_uring fd cache.
|
||||||
|
self.io_backend.invalidate_cached_fds_under(volume, path);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
#[async_recursion::async_recursion]
|
#[async_recursion::async_recursion]
|
||||||
async fn delete_file(
|
async fn delete_file(
|
||||||
@@ -6216,27 +6254,7 @@ impl DiskAPI for LocalDisk {
|
|||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> {
|
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> {
|
||||||
crate::hp_guard!("LocalDisk::delete");
|
crate::hp_guard!("LocalDisk::delete");
|
||||||
let volume_dir = self.get_bucket_path(volume)?;
|
self.delete_unleased(volume, path, &opt).await
|
||||||
if !skip_access_checks(volume)
|
|
||||||
&& let Err(e) = access(&volume_dir).await
|
|
||||||
{
|
|
||||||
return Err(to_access_error(e, DiskError::VolumeAccessDenied).into());
|
|
||||||
}
|
|
||||||
|
|
||||||
let file_path = self.get_object_path(volume, path)?;
|
|
||||||
|
|
||||||
check_path_length(file_path.to_string_lossy().to_string().as_str())?;
|
|
||||||
|
|
||||||
self.delete_file(&volume_dir, &file_path, opt.recursive, opt.immediate)
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
// The inode is unlinked, but a cached descriptor would keep it readable —
|
|
||||||
// a deleted shard must not keep answering reads (backlog#1145). The part
|
|
||||||
// numbers under `path` are not known here, so this is the one caller that
|
|
||||||
// needs the predicate form.
|
|
||||||
self.io_backend.invalidate_cached_fds_under(volume, path);
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
@@ -7824,6 +7842,118 @@ impl DiskAPI for LocalDisk {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result<SnapshotLeaseToken> {
|
||||||
|
let file_path = self.get_object_path(volume, path)?;
|
||||||
|
let key = SnapshotLeaseKey {
|
||||||
|
volume: volume.to_string(),
|
||||||
|
path: path.to_string(),
|
||||||
|
};
|
||||||
|
let token = {
|
||||||
|
let mut registry = self.snapshot_leases.lock().await;
|
||||||
|
if registry.entries.get(&key).is_some_and(|entry| entry.deleting) {
|
||||||
|
return Err(DiskError::FileNotFound);
|
||||||
|
}
|
||||||
|
let token = SnapshotLeaseToken::new();
|
||||||
|
registry.entries.entry(key).or_default().tokens.insert(token);
|
||||||
|
token
|
||||||
|
};
|
||||||
|
match fs::metadata(file_path).await {
|
||||||
|
Ok(metadata) if metadata.is_dir() => Ok(token),
|
||||||
|
Ok(_) => {
|
||||||
|
self.release_snapshot_lease(volume, path, token).await?;
|
||||||
|
Err(DiskError::FileNotFound)
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
self.release_snapshot_lease(volume, path, token).await?;
|
||||||
|
Err(to_file_error(err).into())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> {
|
||||||
|
let key = SnapshotLeaseKey {
|
||||||
|
volume: volume.to_string(),
|
||||||
|
path: path.to_string(),
|
||||||
|
};
|
||||||
|
let opts = {
|
||||||
|
let mut registry = self.snapshot_leases.lock().await;
|
||||||
|
let Some(entry) = registry.entries.get_mut(&key) else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
entry.tokens.remove(&token);
|
||||||
|
if !entry.tokens.is_empty() || entry.deleting {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some(opts) = entry.pending_delete.clone() else {
|
||||||
|
registry.entries.remove(&key);
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
entry.deleting = true;
|
||||||
|
opts
|
||||||
|
};
|
||||||
|
let result = self.delete_unleased(volume, path, &opts).await;
|
||||||
|
let mut registry = self.snapshot_leases.lock().await;
|
||||||
|
match result {
|
||||||
|
Ok(()) => {
|
||||||
|
registry.entries.remove(&key);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
if let Some(entry) = registry.entries.get_mut(&key) {
|
||||||
|
entry.deleting = false;
|
||||||
|
}
|
||||||
|
Err(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result<DataDirDeleteStatus> {
|
||||||
|
let key = SnapshotLeaseKey {
|
||||||
|
volume: volume.to_string(),
|
||||||
|
path: path.to_string(),
|
||||||
|
};
|
||||||
|
{
|
||||||
|
let mut registry = self.snapshot_leases.lock().await;
|
||||||
|
if let Some(entry) = registry.entries.get_mut(&key) {
|
||||||
|
if !entry.tokens.is_empty() {
|
||||||
|
entry.pending_delete.get_or_insert_with(|| opts.clone());
|
||||||
|
return Ok(DataDirDeleteStatus::Deferred);
|
||||||
|
}
|
||||||
|
if entry.deleting {
|
||||||
|
entry.pending_delete.get_or_insert_with(|| opts.clone());
|
||||||
|
return Ok(DataDirDeleteStatus::Deferred);
|
||||||
|
}
|
||||||
|
entry.deleting = true;
|
||||||
|
entry.pending_delete.get_or_insert_with(|| opts.clone());
|
||||||
|
} else {
|
||||||
|
registry.entries.insert(
|
||||||
|
key.clone(),
|
||||||
|
SnapshotLeaseEntry {
|
||||||
|
pending_delete: Some(opts.clone()),
|
||||||
|
deleting: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let result = self.delete_unleased(volume, path, &opts).await;
|
||||||
|
let mut registry = self.snapshot_leases.lock().await;
|
||||||
|
match result {
|
||||||
|
Ok(()) => {
|
||||||
|
registry.entries.remove(&key);
|
||||||
|
Ok(DataDirDeleteStatus::Deleted)
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
if let Some(entry) = registry.entries.get_mut(&key) {
|
||||||
|
entry.deleting = false;
|
||||||
|
}
|
||||||
|
Err(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> {
|
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> {
|
||||||
if !fi.metadata.is_empty() {
|
if !fi.metadata.is_empty() {
|
||||||
@@ -14800,6 +14930,162 @@ mod test {
|
|||||||
assert!(construction_source.is::<ErasureConstructionError>());
|
assert!(construction_source.is::<ErasureConstructionError>());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn snapshot_leases_defer_data_dir_cleanup_until_last_release() {
|
||||||
|
use tempfile::tempdir;
|
||||||
|
|
||||||
|
let root_dir = tempdir().expect("temp dir should be created");
|
||||||
|
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
|
||||||
|
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
|
||||||
|
let volume = "snapshot-lease-volume";
|
||||||
|
let data_dir = path_join_buf(&["object", &Uuid::new_v4().to_string()]);
|
||||||
|
let first_part = path_join_buf(&[&data_dir, "part.1"]);
|
||||||
|
let later_part = path_join_buf(&[&data_dir, "part.2"]);
|
||||||
|
ensure_test_volume(&disk, volume).await;
|
||||||
|
disk.write_all(volume, &first_part, Bytes::from_static(b"first"))
|
||||||
|
.await
|
||||||
|
.expect("first shard should be written");
|
||||||
|
disk.write_all(volume, &later_part, Bytes::from_static(b"later"))
|
||||||
|
.await
|
||||||
|
.expect("later shard should be written");
|
||||||
|
|
||||||
|
let first = disk
|
||||||
|
.acquire_snapshot_lease(volume, &data_dir)
|
||||||
|
.await
|
||||||
|
.expect("first lease should be acquired");
|
||||||
|
let second = disk
|
||||||
|
.acquire_snapshot_lease(volume, &data_dir)
|
||||||
|
.await
|
||||||
|
.expect("second lease should be acquired");
|
||||||
|
let status = disk
|
||||||
|
.delete_data_dir(
|
||||||
|
volume,
|
||||||
|
&data_dir,
|
||||||
|
DeleteOptions {
|
||||||
|
recursive: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("cleanup should be deferred");
|
||||||
|
assert_eq!(status, DataDirDeleteStatus::Deferred);
|
||||||
|
assert_eq!(
|
||||||
|
disk.read_all(volume, &later_part)
|
||||||
|
.await
|
||||||
|
.expect("a later multipart shard must remain openable while leased"),
|
||||||
|
Bytes::from_static(b"later")
|
||||||
|
);
|
||||||
|
|
||||||
|
disk.release_snapshot_lease(volume, &data_dir, first)
|
||||||
|
.await
|
||||||
|
.expect("first lease release should succeed");
|
||||||
|
assert!(
|
||||||
|
disk.read_all(volume, &first_part).await.is_ok(),
|
||||||
|
"one remaining lease must keep the data directory"
|
||||||
|
);
|
||||||
|
disk.release_snapshot_lease(volume, &data_dir, second)
|
||||||
|
.await
|
||||||
|
.expect("last lease release should run deferred cleanup");
|
||||||
|
disk.release_snapshot_lease(volume, &data_dir, second)
|
||||||
|
.await
|
||||||
|
.expect("releasing an already released token should be idempotent");
|
||||||
|
assert!(matches!(disk.read_all(volume, &first_part).await, Err(DiskError::FileNotFound)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn data_dir_cleanup_without_a_lease_keeps_existing_behavior() {
|
||||||
|
use tempfile::tempdir;
|
||||||
|
|
||||||
|
let root_dir = tempdir().expect("temp dir should be created");
|
||||||
|
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
|
||||||
|
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
|
||||||
|
let volume = "snapshot-no-lease-volume";
|
||||||
|
let data_dir = path_join_buf(&["object", &Uuid::new_v4().to_string()]);
|
||||||
|
let part = path_join_buf(&[&data_dir, "part.1"]);
|
||||||
|
ensure_test_volume(&disk, volume).await;
|
||||||
|
disk.write_all(volume, &part, Bytes::from_static(b"payload"))
|
||||||
|
.await
|
||||||
|
.expect("test shard should be written");
|
||||||
|
|
||||||
|
let status = disk
|
||||||
|
.delete_data_dir(
|
||||||
|
volume,
|
||||||
|
&data_dir,
|
||||||
|
DeleteOptions {
|
||||||
|
recursive: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("unleased cleanup should retain the existing delete behavior");
|
||||||
|
assert_eq!(status, DataDirDeleteStatus::Deleted);
|
||||||
|
assert!(matches!(disk.read_all(volume, &part).await, Err(DiskError::FileNotFound)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn snapshot_lease_acquire_and_cleanup_are_atomic() {
|
||||||
|
use tempfile::tempdir;
|
||||||
|
|
||||||
|
let root_dir = tempdir().expect("temp dir should be created");
|
||||||
|
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
|
||||||
|
let disk = Arc::new(LocalDisk::new(&endpoint, false).await.expect("local disk should be created"));
|
||||||
|
let volume = "snapshot-race-volume";
|
||||||
|
ensure_test_volume(&disk, volume).await;
|
||||||
|
|
||||||
|
for iteration in 0..32 {
|
||||||
|
let data_dir = path_join_buf(&["object", &format!("{iteration:032x}")]);
|
||||||
|
let part = path_join_buf(&[&data_dir, "part.1"]);
|
||||||
|
disk.write_all(volume, &part, Bytes::from_static(b"payload"))
|
||||||
|
.await
|
||||||
|
.expect("test shard should be written");
|
||||||
|
let barrier = Arc::new(tokio::sync::Barrier::new(3));
|
||||||
|
let acquire_disk = Arc::clone(&disk);
|
||||||
|
let acquire_barrier = Arc::clone(&barrier);
|
||||||
|
let acquire_path = data_dir.clone();
|
||||||
|
let acquire = tokio::spawn(async move {
|
||||||
|
acquire_barrier.wait().await;
|
||||||
|
acquire_disk.acquire_snapshot_lease(volume, &acquire_path).await
|
||||||
|
});
|
||||||
|
let delete_disk = Arc::clone(&disk);
|
||||||
|
let delete_barrier = Arc::clone(&barrier);
|
||||||
|
let delete_path = data_dir.clone();
|
||||||
|
let delete = tokio::spawn(async move {
|
||||||
|
delete_barrier.wait().await;
|
||||||
|
delete_disk
|
||||||
|
.delete_data_dir(
|
||||||
|
volume,
|
||||||
|
&delete_path,
|
||||||
|
DeleteOptions {
|
||||||
|
recursive: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
barrier.wait().await;
|
||||||
|
|
||||||
|
let acquired = acquire.await.expect("acquire task should join");
|
||||||
|
let deleted = delete
|
||||||
|
.await
|
||||||
|
.expect("delete task should join")
|
||||||
|
.expect("delete should either run or defer");
|
||||||
|
match acquired {
|
||||||
|
Ok(token) => {
|
||||||
|
assert_eq!(deleted, DataDirDeleteStatus::Deferred);
|
||||||
|
assert!(disk.read_all(volume, &part).await.is_ok());
|
||||||
|
disk.release_snapshot_lease(volume, &data_dir, token)
|
||||||
|
.await
|
||||||
|
.expect("release should finish deferred cleanup");
|
||||||
|
}
|
||||||
|
Err(DiskError::FileNotFound) => {
|
||||||
|
assert_eq!(deleted, DataDirDeleteStatus::Deleted);
|
||||||
|
}
|
||||||
|
Err(err) => panic!("unexpected lease acquisition error: {err}"),
|
||||||
|
}
|
||||||
|
assert!(matches!(disk.read_all(volume, &part).await, Err(DiskError::FileNotFound)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn local_disk_check_parts_rejects_zero_data_geometry_before_shard_math() {
|
async fn local_disk_check_parts_rejects_zero_data_geometry_before_shard_math() {
|
||||||
use tempfile::tempdir;
|
use tempfile::tempdir;
|
||||||
|
|||||||
@@ -72,6 +72,27 @@ pub type DiskStore = Arc<Disk>;
|
|||||||
pub type FileReader = Box<dyn AsyncRead + Send + Sync + Unpin>;
|
pub type FileReader = Box<dyn AsyncRead + Send + Sync + Unpin>;
|
||||||
pub type FileWriter = Box<dyn AsyncWrite + Send + Sync + Unpin>;
|
pub type FileWriter = Box<dyn AsyncWrite + Send + Sync + Unpin>;
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize, Deserialize)]
|
||||||
|
pub struct SnapshotLeaseToken(Uuid);
|
||||||
|
|
||||||
|
impl SnapshotLeaseToken {
|
||||||
|
pub fn new() -> Self {
|
||||||
|
Self(Uuid::new_v4())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for SnapshotLeaseToken {
|
||||||
|
fn default() -> Self {
|
||||||
|
Self::new()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
|
pub enum DataDirDeleteStatus {
|
||||||
|
Deleted,
|
||||||
|
Deferred,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
pub enum PartTransactionAction {
|
pub enum PartTransactionAction {
|
||||||
Commit,
|
Commit,
|
||||||
@@ -249,6 +270,27 @@ impl DiskAPI for Disk {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result<SnapshotLeaseToken> {
|
||||||
|
match self {
|
||||||
|
Disk::Local(local_disk) => local_disk.acquire_snapshot_lease(volume, path).await,
|
||||||
|
Disk::Remote(remote_disk) => remote_disk.acquire_snapshot_lease(volume, path).await,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> {
|
||||||
|
match self {
|
||||||
|
Disk::Local(local_disk) => local_disk.release_snapshot_lease(volume, path, token).await,
|
||||||
|
Disk::Remote(remote_disk) => remote_disk.release_snapshot_lease(volume, path, token).await,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result<DataDirDeleteStatus> {
|
||||||
|
match self {
|
||||||
|
Disk::Local(local_disk) => local_disk.delete_data_dir(volume, path, opts).await,
|
||||||
|
Disk::Remote(remote_disk) => remote_disk.delete_data_dir(volume, path, opts).await,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
|
async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
|
||||||
match self {
|
match self {
|
||||||
@@ -646,6 +688,16 @@ pub trait DiskAPI: Debug + Send + Sync + 'static {
|
|||||||
) -> Result<()>;
|
) -> Result<()>;
|
||||||
async fn delete_versions(&self, volume: &str, versions: Vec<FileInfoVersions>, opts: DeleteOptions) -> Vec<Option<Error>>;
|
async fn delete_versions(&self, volume: &str, versions: Vec<FileInfoVersions>, opts: DeleteOptions) -> Vec<Option<Error>>;
|
||||||
async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()>;
|
async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()>;
|
||||||
|
async fn acquire_snapshot_lease(&self, _volume: &str, _path: &str) -> Result<SnapshotLeaseToken> {
|
||||||
|
Err(Error::other("snapshot leases are not supported by this disk"))
|
||||||
|
}
|
||||||
|
async fn release_snapshot_lease(&self, _volume: &str, _path: &str, _token: SnapshotLeaseToken) -> Result<()> {
|
||||||
|
Err(Error::other("snapshot leases are not supported by this disk"))
|
||||||
|
}
|
||||||
|
async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result<DataDirDeleteStatus> {
|
||||||
|
self.delete(volume, path, opts).await?;
|
||||||
|
Ok(DataDirDeleteStatus::Deleted)
|
||||||
|
}
|
||||||
async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()>;
|
async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()>;
|
||||||
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()>;
|
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()>;
|
||||||
async fn read_version(
|
async fn read_version(
|
||||||
|
|||||||
@@ -818,7 +818,11 @@ impl From<s3s::xml::DeError> for Error {
|
|||||||
|
|
||||||
impl From<tonic::Status> for Error {
|
impl From<tonic::Status> for Error {
|
||||||
fn from(e: tonic::Status) -> Self {
|
fn from(e: tonic::Status) -> Self {
|
||||||
Error::other(e.to_string())
|
// Keep the typed status as the io::Error payload instead of a
|
||||||
|
// flattened string so RPC failure classifiers can read the gRPC code
|
||||||
|
// via downcast (see `PeerRestClient::is_network_like_error`).
|
||||||
|
// `RpcStatusError` renders exactly as `e.to_string()` did.
|
||||||
|
Error::other(crate::disk::error::RpcStatusError::from(e))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -180,6 +180,7 @@ pub struct ObjectInfo {
|
|||||||
pub data_dir: Option<Uuid>,
|
pub data_dir: Option<Uuid>,
|
||||||
pub delete_marker: bool,
|
pub delete_marker: bool,
|
||||||
pub transitioned_object: TransitionedObject,
|
pub transitioned_object: TransitionedObject,
|
||||||
|
pub transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
pub restore_ongoing: bool,
|
pub restore_ongoing: bool,
|
||||||
pub restore_expires: Option<OffsetDateTime>,
|
pub restore_expires: Option<OffsetDateTime>,
|
||||||
pub user_tags: Arc<String>,
|
pub user_tags: Arc<String>,
|
||||||
@@ -220,6 +221,7 @@ impl Clone for ObjectInfo {
|
|||||||
data_dir: self.data_dir,
|
data_dir: self.data_dir,
|
||||||
delete_marker: self.delete_marker,
|
delete_marker: self.delete_marker,
|
||||||
transitioned_object: self.transitioned_object.clone(),
|
transitioned_object: self.transitioned_object.clone(),
|
||||||
|
transition_version_state: self.transition_version_state,
|
||||||
restore_ongoing: self.restore_ongoing,
|
restore_ongoing: self.restore_ongoing,
|
||||||
restore_expires: self.restore_expires,
|
restore_expires: self.restore_expires,
|
||||||
user_tags: self.user_tags.clone(),
|
user_tags: self.user_tags.clone(),
|
||||||
@@ -464,11 +466,11 @@ impl ObjectInfo {
|
|||||||
|
|
||||||
let transitioned_object = TransitionedObject {
|
let transitioned_object = TransitionedObject {
|
||||||
name: fi.transitioned_objname.clone(),
|
name: fi.transitioned_objname.clone(),
|
||||||
version_id: if let Some(transition_version_id) = fi.transition_version_id {
|
version_id: fi
|
||||||
transition_version_id.to_string()
|
.transition_version
|
||||||
} else {
|
.clone()
|
||||||
"".to_string()
|
.or_else(|| fi.transition_version_id.map(|version_id| version_id.to_string()))
|
||||||
},
|
.unwrap_or_default(),
|
||||||
status: fi.transition_status.clone(),
|
status: fi.transition_status.clone(),
|
||||||
free_version: fi.tier_free_version(),
|
free_version: fi.tier_free_version(),
|
||||||
tier: fi.transition_tier.clone(),
|
tier: fi.transition_tier.clone(),
|
||||||
@@ -537,6 +539,7 @@ impl ObjectInfo {
|
|||||||
inlined,
|
inlined,
|
||||||
user_defined: Arc::new(metadata),
|
user_defined: Arc::new(metadata),
|
||||||
transitioned_object,
|
transitioned_object,
|
||||||
|
transition_version_state: fi.transition_version_state,
|
||||||
checksum: fi.checksum.clone(),
|
checksum: fi.checksum.clone(),
|
||||||
storage_class,
|
storage_class,
|
||||||
restore_ongoing,
|
restore_ongoing,
|
||||||
|
|||||||
@@ -931,7 +931,10 @@ pub async fn read_transition_meta(disk_path: &Path, bucket: &str, object: &str)
|
|||||||
status: fi.transition_status.clone(),
|
status: fi.transition_status.clone(),
|
||||||
tier: fi.transition_tier.clone(),
|
tier: fi.transition_tier.clone(),
|
||||||
remote_object: fi.transitioned_objname.clone(),
|
remote_object: fi.transitioned_objname.clone(),
|
||||||
remote_version_id: fi.transition_version_id.map(|id| id.to_string()),
|
remote_version_id: fi
|
||||||
|
.transition_version
|
||||||
|
.clone()
|
||||||
|
.or_else(|| fi.transition_version_id.map(|id| id.to_string())),
|
||||||
free_version_count,
|
free_version_count,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9274,7 +9274,8 @@ mod tests {
|
|||||||
version_id: "v1".to_string(),
|
version_id: "v1".to_string(),
|
||||||
tier_name: "COLD-A".to_string(),
|
tier_name: "COLD-A".to_string(),
|
||||||
backend_identity: Some(current_identity),
|
backend_identity: Some(current_identity),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
};
|
};
|
||||||
journal_store
|
journal_store
|
||||||
.insert_config_object(
|
.insert_config_object(
|
||||||
|
|||||||
@@ -47,8 +47,8 @@ use crate::diagnostics::get::{
|
|||||||
record_get_object_pipeline_failure_for_path, record_get_stage_duration_if_enabled,
|
record_get_object_pipeline_failure_for_path, record_get_stage_duration_if_enabled,
|
||||||
};
|
};
|
||||||
use crate::disk::{
|
use crate::disk::{
|
||||||
OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction,
|
DataDirDeleteStatus, OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK,
|
||||||
part_transaction_path,
|
PartTransactionAction, part_transaction_path,
|
||||||
};
|
};
|
||||||
use crate::erasure::coding::BitrotReader;
|
use crate::erasure::coding::BitrotReader;
|
||||||
use crate::io_support::bitrot::ShardReader;
|
use crate::io_support::bitrot::ShardReader;
|
||||||
@@ -547,6 +547,8 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo
|
|||||||
&& left.transitioned_objname == right.transitioned_objname
|
&& left.transitioned_objname == right.transitioned_objname
|
||||||
&& left.transition_tier == right.transition_tier
|
&& left.transition_tier == right.transition_tier
|
||||||
&& left.transition_version_id == right.transition_version_id
|
&& left.transition_version_id == right.transition_version_id
|
||||||
|
&& left.transition_version == right.transition_version
|
||||||
|
&& left.transition_version_state == right.transition_version_state
|
||||||
&& left.expire_restored == right.expire_restored
|
&& left.expire_restored == right.expire_restored
|
||||||
&& left.size == right.size
|
&& left.size == right.size
|
||||||
&& left.mod_time == right.mod_time
|
&& left.mod_time == right.mod_time
|
||||||
@@ -2985,29 +2987,48 @@ impl SetDisks {
|
|||||||
Self::rename_fanout_barrier(&object_for_fault, idx, rename_fanout_barrier_phase::CLEANUP).await;
|
Self::rename_fanout_barrier(&object_for_fault, idx, rename_fanout_barrier_phase::CLEANUP).await;
|
||||||
|
|
||||||
if let Some(err) = Self::cleanup_injected_error(&object_for_fault, idx) {
|
if let Some(err) = Self::cleanup_injected_error(&object_for_fault, idx) {
|
||||||
return Some(err);
|
return (false, Some(err));
|
||||||
}
|
}
|
||||||
if let Some(disk) = disk {
|
if let Some(disk) = disk {
|
||||||
disk.delete(
|
match disk
|
||||||
&bucket,
|
.delete_data_dir(
|
||||||
&file_path,
|
&bucket,
|
||||||
DeleteOptions {
|
&file_path,
|
||||||
recursive: true,
|
DeleteOptions {
|
||||||
..Default::default()
|
recursive: true,
|
||||||
},
|
..Default::default()
|
||||||
)
|
},
|
||||||
.await
|
)
|
||||||
.err()
|
.await
|
||||||
|
{
|
||||||
|
Ok(DataDirDeleteStatus::Deleted) => (false, None),
|
||||||
|
Ok(DataDirDeleteStatus::Deferred) => (true, None),
|
||||||
|
Err(err) => (false, Some(err)),
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
// `None` slot: ignored placeholder. It is not `attempted`, so
|
// `None` slot: ignored placeholder. It is not `attempted`, so
|
||||||
// classification excludes it from residue regardless.
|
// classification excludes it from residue regardless.
|
||||||
Some(DiskError::DiskNotFound)
|
(false, Some(DiskError::DiskNotFound))
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
});
|
});
|
||||||
let errs: Vec<Option<DiskError>> = join_all(futures).await.into_iter().map(map_cleanup_join_result).collect();
|
let mut deferred = 0usize;
|
||||||
|
let errs: Vec<Option<DiskError>> = join_all(futures)
|
||||||
|
.await
|
||||||
|
.into_iter()
|
||||||
|
.map(|result| match result {
|
||||||
|
Ok((was_deferred, err)) => {
|
||||||
|
deferred += usize::from(was_deferred);
|
||||||
|
err
|
||||||
|
}
|
||||||
|
Err(join_err) => Some(DiskError::other(format!("old data dir cleanup task failed: {join_err}"))),
|
||||||
|
})
|
||||||
|
.collect();
|
||||||
|
|
||||||
classify_old_data_dir_cleanup(&errs, &attempted, write_quorum)
|
let mut cleanup = classify_old_data_dir_cleanup(&errs, &attempted, write_quorum);
|
||||||
|
cleanup.deferred = deferred;
|
||||||
|
cleanup.reclaimed = cleanup.reclaimed.saturating_sub(deferred);
|
||||||
|
cleanup
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Test-only fault-injection seam for the old-data-dir cleanup path
|
/// Test-only fault-injection seam for the old-data-dir cleanup path
|
||||||
@@ -3098,6 +3119,20 @@ impl SetDisks {
|
|||||||
|
|
||||||
rustfs_io_metrics::record_old_data_dir_cleanup(c.attempted, c.reclaimed, c.unreclaimed_disks.len(), c.below_quorum);
|
rustfs_io_metrics::record_old_data_dir_cleanup(c.attempted, c.reclaimed, c.unreclaimed_disks.len(), c.below_quorum);
|
||||||
|
|
||||||
|
if c.deferred > 0 {
|
||||||
|
debug!(
|
||||||
|
event = EVENT_SET_DISK_WRITE,
|
||||||
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
|
subsystem = LOG_SUBSYSTEM_SET_DISK,
|
||||||
|
bucket = %bucket,
|
||||||
|
object = %object,
|
||||||
|
old_data_dir = %old_dir,
|
||||||
|
deferred = c.deferred,
|
||||||
|
state = "old_data_cleanup_deferred",
|
||||||
|
"Old data directory cleanup deferred for active snapshot leases"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
if actions.warn {
|
if actions.warn {
|
||||||
warn!(
|
warn!(
|
||||||
component = LOG_COMPONENT_ECSTORE,
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
@@ -4129,6 +4164,9 @@ pub(in crate::set_disk) struct OldDataDirCleanup {
|
|||||||
/// Number of attempted disks that returned `Ok` or a not-found variant
|
/// Number of attempted disks that returned `Ok` or a not-found variant
|
||||||
/// (a missing dir == already reclaimed).
|
/// (a missing dir == already reclaimed).
|
||||||
pub reclaimed: usize,
|
pub reclaimed: usize,
|
||||||
|
/// Number of attempted disks that retained the directory for an active
|
||||||
|
/// snapshot lease and registered it for deletion after the final release.
|
||||||
|
pub deferred: usize,
|
||||||
/// Indices of attempted disks that failed with a non-ignored, non-not-found
|
/// Indices of attempted disks that failed with a non-ignored, non-not-found
|
||||||
/// error (including task panic/cancel). This is the residue that actually
|
/// error (including task panic/cancel). This is the residue that actually
|
||||||
/// leaks and drives the leak metric + heal enqueue.
|
/// leaks and drives the leak metric + heal enqueue.
|
||||||
@@ -4191,6 +4229,7 @@ fn classify_old_data_dir_cleanup(errs: &[Option<DiskError>], attempted: &[bool],
|
|||||||
OldDataDirCleanup {
|
OldDataDirCleanup {
|
||||||
attempted: attempted_count,
|
attempted: attempted_count,
|
||||||
reclaimed,
|
reclaimed,
|
||||||
|
deferred: 0,
|
||||||
unreclaimed_disks,
|
unreclaimed_disks,
|
||||||
below_quorum,
|
below_quorum,
|
||||||
}
|
}
|
||||||
@@ -5131,6 +5170,40 @@ mod tests {
|
|||||||
drop((disk1, disk2));
|
drop((disk1, disk2));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn commit_cleanup_reports_and_releases_deferred_snapshot_data_dirs() {
|
||||||
|
let bucket = "cleanup-lease-bucket";
|
||||||
|
let object = "cleanup-lease-object";
|
||||||
|
let old_data_dir = "11111111-1111-1111-1111-111111111111";
|
||||||
|
let committed_data_dir = "22222222-2222-2222-2222-222222222222";
|
||||||
|
let data_dir_path = format!("{object}/{old_data_dir}");
|
||||||
|
let shard_path = format!("{data_dir_path}/part.1");
|
||||||
|
let (_dir1, disk1) = read_multiple_test_disk(bucket, &[(&shard_path, b"one".as_slice())]).await;
|
||||||
|
let set = io_primitives_test_set(vec![Some(disk1.clone())], 0).await;
|
||||||
|
let lease = disk1
|
||||||
|
.acquire_snapshot_lease(bucket, &data_dir_path)
|
||||||
|
.await
|
||||||
|
.expect("snapshot lease should be acquired before cleanup");
|
||||||
|
|
||||||
|
let cleanup = set
|
||||||
|
.commit_rename_data_dir(&[Some(disk1.clone())], bucket, object, old_data_dir, committed_data_dir, 1)
|
||||||
|
.await;
|
||||||
|
assert_eq!(cleanup.attempted, 1);
|
||||||
|
assert_eq!(cleanup.reclaimed, 0);
|
||||||
|
assert_eq!(cleanup.deferred, 1);
|
||||||
|
assert!(cleanup.unreclaimed_disks.is_empty());
|
||||||
|
disk1
|
||||||
|
.read_all(bucket, &shard_path)
|
||||||
|
.await
|
||||||
|
.expect("deferred cleanup must leave later shard opens available");
|
||||||
|
|
||||||
|
disk1
|
||||||
|
.release_snapshot_lease(bucket, &data_dir_path, lease)
|
||||||
|
.await
|
||||||
|
.expect("final lease release should reclaim the old data directory");
|
||||||
|
assert!(matches!(disk1.read_all(bucket, &shard_path).await, Err(DiskError::FileNotFound)));
|
||||||
|
}
|
||||||
|
|
||||||
/// Isolation guard: an armed barrier / observed object only affects its own
|
/// Isolation guard: an armed barrier / observed object only affects its own
|
||||||
/// object. A fan-out for a different (unobserved, unarmed) object must not be
|
/// object. A fan-out for a different (unobserved, unarmed) object must not be
|
||||||
/// paused and must not accrue any tracked task count — so concurrent tests
|
/// paused and must not accrue any tracked task count — so concurrent tests
|
||||||
|
|||||||
@@ -578,6 +578,13 @@ impl SetDisks {
|
|||||||
Self::update_hash_str(hasher, &meta.transition_tier);
|
Self::update_hash_str(hasher, &meta.transition_tier);
|
||||||
Self::update_hash_str(hasher, &meta.transitioned_objname);
|
Self::update_hash_str(hasher, &meta.transitioned_objname);
|
||||||
Self::update_hash_optional_uuid(hasher, meta.transition_version_id);
|
Self::update_hash_optional_uuid(hasher, meta.transition_version_id);
|
||||||
|
Self::update_hash_optional_str(hasher, meta.transition_version.as_deref());
|
||||||
|
hasher.update([match meta.transition_version_state {
|
||||||
|
rustfs_filemeta::TransitionVersionState::Unknown => 0,
|
||||||
|
rustfs_filemeta::TransitionVersionState::KnownDisabled => 1,
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull => 2,
|
||||||
|
rustfs_filemeta::TransitionVersionState::Exact => 3,
|
||||||
|
}]);
|
||||||
Self::update_hash_optional_u32(hasher, meta.mode);
|
Self::update_hash_optional_u32(hasher, meta.mode);
|
||||||
Self::update_hash_optional_u64(hasher, meta.written_by_version);
|
Self::update_hash_optional_u64(hasher, meta.written_by_version);
|
||||||
|
|
||||||
|
|||||||
@@ -692,6 +692,8 @@ pub(crate) use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
|||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
||||||
pub(crate) use ops::object::body_cache_plaintext_len;
|
pub(crate) use ops::object::body_cache_plaintext_len;
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) use ops::object::cleanup_rejected_transition_upload_durably;
|
||||||
mod read;
|
mod read;
|
||||||
mod replication;
|
mod replication;
|
||||||
pub(crate) mod shard_source;
|
pub(crate) mod shard_source;
|
||||||
|
|||||||
@@ -25,7 +25,10 @@ use crate::set_disk::read::GetObjectDownstreamWriter;
|
|||||||
|
|
||||||
use crate::bucket::lifecycle::{
|
use crate::bucket::lifecycle::{
|
||||||
tier_delete_journal::{persist_tier_delete_journal_entry, remove_tier_delete_journal_entry},
|
tier_delete_journal::{persist_tier_delete_journal_entry, remove_tier_delete_journal_entry},
|
||||||
tier_sweeper::{Jentry, RemoteTierDeleteOutcome, delete_object_from_remote_tier_with_lease_idempotent},
|
tier_sweeper::{
|
||||||
|
Jentry, RemoteTierDeleteOutcome, delete_confirmed_transition_candidate_exact_with_lease_idempotent,
|
||||||
|
delete_object_from_remote_tier_with_lease_idempotent,
|
||||||
|
},
|
||||||
transition_transaction::{
|
transition_transaction::{
|
||||||
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
||||||
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
|
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
|
||||||
@@ -1663,7 +1666,11 @@ pub(crate) async fn cleanup_uncommitted_transition_upload(
|
|||||||
cleanup_version: &str,
|
cleanup_version: &str,
|
||||||
version_id_exact: bool,
|
version_id_exact: bool,
|
||||||
) -> std::io::Result<RemoteTierDeleteOutcome> {
|
) -> std::io::Result<RemoteTierDeleteOutcome> {
|
||||||
delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, version_id_exact).await
|
if version_id_exact {
|
||||||
|
delete_confirmed_transition_candidate_exact_with_lease_idempotent(object, cleanup_version, lease).await
|
||||||
|
} else {
|
||||||
|
delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, false).await
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn log_transition_upload_cleanup_failure(lease: &TierOperationLease, object: &str, cleanup_version: &str, err: &std::io::Error) {
|
fn log_transition_upload_cleanup_failure(lease: &TierOperationLease, object: &str, cleanup_version: &str, err: &std::io::Error) {
|
||||||
@@ -1800,7 +1807,7 @@ impl Drop for TransitionUploadCleanup {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn cleanup_rejected_transition_upload_durably(
|
pub(crate) async fn cleanup_rejected_transition_upload_durably(
|
||||||
lease: &TierOperationLease,
|
lease: &TierOperationLease,
|
||||||
object: &str,
|
object: &str,
|
||||||
cleanup_version: &str,
|
cleanup_version: &str,
|
||||||
@@ -1813,6 +1820,13 @@ async fn cleanup_rejected_transition_upload_durably(
|
|||||||
tier_name: lease.tier_name().to_string(),
|
tier_name: lease.tier_name().to_string(),
|
||||||
backend_identity: Some(lease.backend_identity()),
|
backend_identity: Some(lease.backend_identity()),
|
||||||
version_id_exact,
|
version_id_exact,
|
||||||
|
version_state: if !version_id_exact {
|
||||||
|
rustfs_filemeta::TransitionVersionState::KnownDisabled
|
||||||
|
} else if cleanup_version == "null" {
|
||||||
|
rustfs_filemeta::TransitionVersionState::SuspendedNull
|
||||||
|
} else {
|
||||||
|
rustfs_filemeta::TransitionVersionState::Exact
|
||||||
|
},
|
||||||
};
|
};
|
||||||
|
|
||||||
let journal_error = if let Some(api) = api.as_ref() {
|
let journal_error = if let Some(api) = api.as_ref() {
|
||||||
@@ -1972,12 +1986,83 @@ async fn advance_and_save_transition_transaction(
|
|||||||
next: TransitionTransactionState,
|
next: TransitionTransactionState,
|
||||||
remote_version: Option<TransitionRemoteVersion>,
|
remote_version: Option<TransitionRemoteVersion>,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
|
#[cfg(test)]
|
||||||
|
record_transition_uploaded_save_attempt(transaction, next);
|
||||||
transaction
|
transaction
|
||||||
.advance(transaction.fence(), next, remote_version)
|
.advance(transaction.fence(), next, remote_version)
|
||||||
.map_err(Error::other)?;
|
.map_err(Error::other)?;
|
||||||
save_transition_transaction_if_available(api, transaction).await
|
save_transition_transaction_if_available(api, transaction).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
struct TransitionUploadedSaveProbeState {
|
||||||
|
bucket: String,
|
||||||
|
object: String,
|
||||||
|
attempts: std::sync::atomic::AtomicUsize,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
struct TransitionUploadedSaveProbe {
|
||||||
|
state: Arc<TransitionUploadedSaveProbeState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static TRANSITION_UPLOADED_SAVE_PROBE: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TransitionUploadedSaveProbeState>>>> =
|
||||||
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl TransitionUploadedSaveProbe {
|
||||||
|
fn install(bucket: &str, object: &str) -> Self {
|
||||||
|
let state = Arc::new(TransitionUploadedSaveProbeState {
|
||||||
|
bucket: bucket.to_string(),
|
||||||
|
object: object.to_string(),
|
||||||
|
attempts: std::sync::atomic::AtomicUsize::new(0),
|
||||||
|
});
|
||||||
|
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("transition uploaded-save probe mutex should not poison");
|
||||||
|
assert!(slot.is_none(), "transition uploaded-save probe must be installed by one test at a time");
|
||||||
|
*slot = Some(Arc::clone(&state));
|
||||||
|
drop(slot);
|
||||||
|
Self { state }
|
||||||
|
}
|
||||||
|
|
||||||
|
fn attempts(&self) -> usize {
|
||||||
|
self.state.attempts.load(std::sync::atomic::Ordering::Acquire)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
impl Drop for TransitionUploadedSaveProbe {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
let mut slot = TRANSITION_UPLOADED_SAVE_PROBE
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("transition uploaded-save probe mutex should not poison");
|
||||||
|
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
||||||
|
*slot = None;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction, next: TransitionTransactionState) {
|
||||||
|
if next != TransitionTransactionState::Uploaded {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let state = TRANSITION_UPLOADED_SAVE_PROBE
|
||||||
|
.get_or_init(|| std::sync::Mutex::new(None))
|
||||||
|
.lock()
|
||||||
|
.expect("transition uploaded-save probe mutex should not poison")
|
||||||
|
.as_ref()
|
||||||
|
.filter(|state| state.bucket == transaction.source.bucket && state.object == transaction.source.object)
|
||||||
|
.cloned();
|
||||||
|
if let Some(state) = state {
|
||||||
|
state.attempts.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
|
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
|
||||||
if let Some(api) = api {
|
if let Some(api) = api {
|
||||||
return delete_transition_transaction_record(api.clone(), transaction_id).await;
|
return delete_transition_transaction_record(api.clone(), transaction_id).await;
|
||||||
@@ -2228,11 +2313,25 @@ async fn pause_transition_commit(bucket: &str, object: &str, pause: TransitionCo
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn parse_transition_version_id(remote_version: &str) -> std::result::Result<Option<Uuid>, uuid::Error> {
|
fn persisted_transition_version(
|
||||||
|
remote_version: &str,
|
||||||
|
) -> std::io::Result<(Option<String>, rustfs_filemeta::TransitionVersionState)> {
|
||||||
if remote_version.is_empty() {
|
if remote_version.is_empty() {
|
||||||
return Ok(None);
|
return Ok((None, rustfs_filemeta::TransitionVersionState::KnownDisabled));
|
||||||
}
|
}
|
||||||
Uuid::parse_str(remote_version).map(|version_id| (!version_id.is_nil()).then_some(version_id))
|
let version_id = Uuid::parse_str(remote_version).map_err(|_| {
|
||||||
|
std::io::Error::new(
|
||||||
|
std::io::ErrorKind::Unsupported,
|
||||||
|
"opaque remote tier versions require the cluster capability gate",
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
if version_id.is_nil() {
|
||||||
|
return Err(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::InvalidData,
|
||||||
|
"remote tier returned a nil object version ID",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
Ok((Some(remote_version.to_string()), rustfs_filemeta::TransitionVersionState::Exact))
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -2470,16 +2569,17 @@ mod transition_upload_completion_tests {
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod transition_version_id_tests {
|
mod transition_version_id_tests {
|
||||||
use super::{TransitionUploadCandidate, parse_transition_version_id};
|
use super::{TransitionUploadCandidate, persisted_transition_version};
|
||||||
|
use rustfs_filemeta::TransitionVersionState;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() {
|
fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() {
|
||||||
assert_eq!(parse_transition_version_id("").expect("empty remote version should be valid"), None);
|
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
parse_transition_version_id(&Uuid::nil().to_string()).expect("nil remote version should be valid"),
|
persisted_transition_version("").expect("empty remote version identifies an unversioned tier"),
|
||||||
None
|
(None, TransitionVersionState::KnownDisabled)
|
||||||
);
|
);
|
||||||
|
assert!(persisted_transition_version(&Uuid::nil().to_string()).is_err());
|
||||||
let nil_put_response = Uuid::nil().to_string();
|
let nil_put_response = Uuid::nil().to_string();
|
||||||
let nil_candidate = TransitionUploadCandidate::from_put_response(nil_put_response.clone());
|
let nil_candidate = TransitionUploadCandidate::from_put_response(nil_put_response.clone());
|
||||||
assert_eq!(nil_candidate.cleanup_version(), nil_put_response);
|
assert_eq!(nil_candidate.cleanup_version(), nil_put_response);
|
||||||
@@ -2491,12 +2591,14 @@ mod transition_version_id_tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn preserves_valid_remote_id_and_rejects_invalid_text() {
|
fn preserves_uuid_and_gates_opaque_remote_ids() {
|
||||||
let version_id = Uuid::new_v4();
|
let version_id = Uuid::new_v4();
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
parse_transition_version_id(&version_id.to_string()).expect("UUID remote version should be valid"),
|
persisted_transition_version(&version_id.to_string()).expect("UUID remote version"),
|
||||||
Some(version_id)
|
(Some(version_id.to_string()), TransitionVersionState::Exact)
|
||||||
);
|
);
|
||||||
|
assert!(persisted_transition_version("null").is_err());
|
||||||
|
assert!(persisted_transition_version("opaque-version-token").is_err());
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
TransitionUploadCandidate::from_put_response(version_id.to_string()).cleanup_version(),
|
TransitionUploadCandidate::from_put_response(version_id.to_string()).cleanup_version(),
|
||||||
version_id.to_string()
|
version_id.to_string()
|
||||||
@@ -2505,7 +2607,6 @@ mod transition_version_id_tests {
|
|||||||
TransitionUploadCandidate::from_put_response("opaque-version-token".to_string()).cleanup_version(),
|
TransitionUploadCandidate::from_put_response("opaque-version-token".to_string()).cleanup_version(),
|
||||||
"opaque-version-token"
|
"opaque-version-token"
|
||||||
);
|
);
|
||||||
assert!(parse_transition_version_id("not-a-uuid").is_err());
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3775,6 +3876,20 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
|
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
|
||||||
return Err(err.into());
|
return Err(err.into());
|
||||||
}
|
}
|
||||||
|
let (transition_version_id, transition_version_state) = match persisted_transition_version(candidate.remote_version()) {
|
||||||
|
Ok(version) => version,
|
||||||
|
Err(err) => {
|
||||||
|
let cleanup_api = transition_cleanup_store(&self.ctx).await;
|
||||||
|
if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api).await {
|
||||||
|
return Err(StorageError::Io(std::io::Error::other(format!(
|
||||||
|
"{err}; rejected remote upload cleanup failed: {cleanup_err}"
|
||||||
|
))));
|
||||||
|
}
|
||||||
|
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object)
|
||||||
|
.await;
|
||||||
|
return Err(err.into());
|
||||||
|
}
|
||||||
|
};
|
||||||
if let Err(err) = advance_and_save_transition_transaction(
|
if let Err(err) = advance_and_save_transition_transaction(
|
||||||
transaction_api.as_ref(),
|
transaction_api.as_ref(),
|
||||||
&mut transaction,
|
&mut transaction,
|
||||||
@@ -3792,16 +3907,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
|
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await;
|
||||||
return Err(err);
|
return Err(err);
|
||||||
}
|
}
|
||||||
let transition_version_id = match parse_transition_version_id(candidate.remote_version()) {
|
|
||||||
Ok(version_id) => version_id,
|
|
||||||
Err(err) => {
|
|
||||||
if upload_cleanup.cleanup().await.is_ok() {
|
|
||||||
delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object)
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
return Err(err.into());
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
let mut commit_opts = opts.clone();
|
let mut commit_opts = opts.clone();
|
||||||
commit_opts.no_lock = true;
|
commit_opts.no_lock = true;
|
||||||
@@ -3859,7 +3964,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
current_fi.transition_status = TRANSITION_COMPLETE.to_string();
|
current_fi.transition_status = TRANSITION_COMPLETE.to_string();
|
||||||
current_fi.transitioned_objname = dest_obj;
|
current_fi.transitioned_objname = dest_obj;
|
||||||
current_fi.transition_tier = opts.transition.tier.clone();
|
current_fi.transition_tier = opts.transition.tier.clone();
|
||||||
current_fi.transition_version_id = transition_version_id;
|
current_fi.transition_version_id = transition_version_id
|
||||||
|
.as_deref()
|
||||||
|
.and_then(|version_id| Uuid::parse_str(version_id).ok());
|
||||||
|
current_fi.transition_version = transition_version_id;
|
||||||
|
current_fi.transition_version_state = transition_version_state;
|
||||||
rustfs_utils::http::metadata_compat::insert_str(
|
rustfs_utils::http::metadata_compat::insert_str(
|
||||||
&mut current_fi.metadata,
|
&mut current_fi.metadata,
|
||||||
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||||
@@ -4076,6 +4185,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
let mut p_reader = PutObjReader::new(hash_reader);
|
let mut p_reader = PutObjReader::new(hash_reader);
|
||||||
return match self_.clone().put_object(bucket, object, &mut p_reader, &ropts).await {
|
return match self_.clone().put_object(bucket, object, &mut p_reader, &ropts).await {
|
||||||
Ok(restored_info) => {
|
Ok(restored_info) => {
|
||||||
|
let restored_info = self_.finalize_restore_metadata(bucket, object, &restored_info, &opts).await?;
|
||||||
send_event(EventArgs {
|
send_event(EventArgs {
|
||||||
event_name: EventName::ObjectRestoreCompleted.as_str().to_string(),
|
event_name: EventName::ObjectRestoreCompleted.as_str().to_string(),
|
||||||
bucket_name: bucket.to_string(),
|
bucket_name: bucket.to_string(),
|
||||||
@@ -4210,6 +4320,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
return set_restore_header_fn(&mut oi, Some(err)).await;
|
return set_restore_header_fn(&mut oi, Some(err)).await;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
let restored_info = self_.finalize_restore_metadata(bucket, object, &restored_info, opts).await?;
|
||||||
send_event(EventArgs {
|
send_event(EventArgs {
|
||||||
event_name: EventName::ObjectRestoreCompleted.as_str().to_string(),
|
event_name: EventName::ObjectRestoreCompleted.as_str().to_string(),
|
||||||
bucket_name: bucket.to_string(),
|
bucket_name: bucket.to_string(),
|
||||||
@@ -4668,6 +4779,33 @@ mod transition_commit_failure_tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
async fn rejected_unsupported_remote_versions_are_cleaned_up() {
|
||||||
|
for remote_version in ["null", "opaque-version-token"] {
|
||||||
|
let manager = TierConfigMgr::new();
|
||||||
|
let backend = register_mock_tier(&manager, "WARM").await;
|
||||||
|
let lease = TierConfigMgr::acquire_operation_lease(&manager, "WARM")
|
||||||
|
.await
|
||||||
|
.expect("mock tier lease should be available");
|
||||||
|
let candidate = TransitionUploadCandidate::from_put_response(remote_version.to_string());
|
||||||
|
|
||||||
|
persisted_transition_version(candidate.remote_version()).expect_err("unsupported writer version must fail closed");
|
||||||
|
cleanup_rejected_transition_upload_durably(
|
||||||
|
&lease,
|
||||||
|
"remote/object",
|
||||||
|
candidate.cleanup_version(),
|
||||||
|
candidate.cleanup_version_is_exact(),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("rejected remote upload must be cleaned up");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
backend.remove_versions().await,
|
||||||
|
vec![("remote/object".to_string(), candidate.cleanup_version().to_string())]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[serial_test::serial(restore_multipart_failure_point)]
|
#[serial_test::serial(restore_multipart_failure_point)]
|
||||||
async fn multipart_restore_aborts_every_post_create_failure() {
|
async fn multipart_restore_aborts_every_post_create_failure() {
|
||||||
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||||
@@ -6615,6 +6753,45 @@ mod transition_upload_integrity_tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn unversioned_remote_version_is_persisted_without_version_id() {
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||||
|
let bucket = "transition-unversioned-tier-bucket";
|
||||||
|
let object = "object.bin";
|
||||||
|
let payload = b"unversioned remote tier must commit without a version id".repeat(1024);
|
||||||
|
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
|
||||||
|
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
|
||||||
|
backend.set_put_remote_version(Some(String::new())).await;
|
||||||
|
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
|
||||||
|
|
||||||
|
set_disks
|
||||||
|
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
||||||
|
.await
|
||||||
|
.expect("an unversioned remote version must commit");
|
||||||
|
let (fi, _, _) = set_disks
|
||||||
|
.get_object_fileinfo(
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
&ObjectOptions {
|
||||||
|
no_lock: true,
|
||||||
|
metadata_cache_safe: false,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("committed unversioned transition metadata should be readable");
|
||||||
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(fi.transition_version, None);
|
||||||
|
assert_eq!(fi.transition_version_state, rustfs_filemeta::TransitionVersionState::KnownDisabled);
|
||||||
|
assert_eq!(save_probe.attempts(), 1);
|
||||||
|
assert_eq!(backend.remove_count().await, 0);
|
||||||
|
assert_eq!(backend.object_count().await, 1);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn opaque_remote_version_is_cleaned_before_parse_failure() {
|
async fn opaque_remote_version_is_cleaned_before_parse_failure() {
|
||||||
@@ -6630,7 +6807,7 @@ mod transition_upload_integrity_tests {
|
|||||||
set_disks
|
set_disks
|
||||||
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
||||||
.await
|
.await
|
||||||
.expect_err("an unparseable remote version must fail closed");
|
.expect_err("an opaque remote version must fail closed until the capability gate is active");
|
||||||
let removed_versions = backend.remove_versions().await;
|
let removed_versions = backend.remove_versions().await;
|
||||||
assert_eq!(removed_versions.len(), 1);
|
assert_eq!(removed_versions.len(), 1);
|
||||||
assert_eq!(removed_versions[0].1, "opaque-version-token");
|
assert_eq!(removed_versions[0].1, "opaque-version-token");
|
||||||
@@ -6638,6 +6815,38 @@ mod transition_upload_integrity_tests {
|
|||||||
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn nil_remote_version_is_cleaned_exactly_before_transaction_persistence() {
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||||
|
let bucket = "transition-nil-version-bucket";
|
||||||
|
let object = "object.bin";
|
||||||
|
let payload = b"nil remote version must retain local data".repeat(1024);
|
||||||
|
let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await;
|
||||||
|
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||||
|
let remote_version = Uuid::nil().to_string();
|
||||||
|
let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await;
|
||||||
|
backend.set_put_remote_version(Some(remote_version.clone())).await;
|
||||||
|
let save_probe = TransitionUploadedSaveProbe::install(bucket, object);
|
||||||
|
|
||||||
|
set_disks
|
||||||
|
.transition_object(bucket, object, &transition_options(&original, tier_name))
|
||||||
|
.await
|
||||||
|
.expect_err("a nil remote version must fail closed before transaction persistence");
|
||||||
|
let put_versions = backend.put_versions().await;
|
||||||
|
let removed_versions = backend.remove_versions().await;
|
||||||
|
assert_eq!(removed_versions, put_versions);
|
||||||
|
assert_eq!(removed_versions.len(), 1);
|
||||||
|
assert_eq!(
|
||||||
|
removed_versions.first().map(|(_, version)| version.as_str()),
|
||||||
|
Some(remote_version.as_str())
|
||||||
|
);
|
||||||
|
assert_eq!(save_probe.attempts(), 0, "nil remote version must be rejected before saving Uploaded");
|
||||||
|
assert_eq!(backend.exact_remove_count(), 1);
|
||||||
|
assert_eq!(backend.object_count().await, 0);
|
||||||
|
assert_local_source_intact(&set_disks, bucket, object, &payload).await;
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn authoritative_read_failure_after_upload_cleans_exact_candidate_and_preserves_source() {
|
async fn authoritative_read_failure_after_upload_cleans_exact_candidate_and_preserves_source() {
|
||||||
@@ -6985,23 +7194,36 @@ mod transition_source_identity_matrix_tests {
|
|||||||
let object = format!("identity-{index}.bin");
|
let object = format!("identity-{index}.bin");
|
||||||
let payload = vec![u8::try_from(index + 1).expect("matrix index should fit u8"); 1024 * 1024];
|
let payload = vec![u8::try_from(index + 1).expect("matrix index should fit u8"); 1024 * 1024];
|
||||||
let mut reader = PutObjReader::from_vec(payload);
|
let mut reader = PutObjReader::from_vec(payload);
|
||||||
|
let source_version_id = Uuid::new_v4();
|
||||||
|
let source_opts = ObjectOptions {
|
||||||
|
version_id: Some(source_version_id.to_string()),
|
||||||
|
versioned: true,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
let original = set_disks
|
let original = set_disks
|
||||||
.put_object(bucket, &object, &mut reader, &ObjectOptions::default())
|
.put_object(bucket, &object, &mut reader, &source_opts)
|
||||||
.await
|
.await
|
||||||
.expect("source object should be written");
|
.expect("source object should be written");
|
||||||
let (source, _, _) = set_disks
|
let (source, _, _) = set_disks
|
||||||
.get_object_fileinfo(bucket, &object, &ObjectOptions::default(), true, false)
|
.get_object_fileinfo(bucket, &object, &source_opts, true, false)
|
||||||
.await
|
.await
|
||||||
.expect("source metadata should resolve");
|
.expect("source metadata should resolve");
|
||||||
|
assert_eq!(source.version_id, Some(source_version_id));
|
||||||
|
assert_eq!(
|
||||||
|
transition_source_identity(bucket, &object, &source, &source_opts, &get_raw_etag(&source.metadata))
|
||||||
|
.expect("persisted versioned source identity should build")
|
||||||
|
.version_mode,
|
||||||
|
TransitionSourceVersionMode::Versioned
|
||||||
|
);
|
||||||
let opts = ObjectOptions {
|
let opts = ObjectOptions {
|
||||||
no_lock: true,
|
no_lock: true,
|
||||||
|
versioned: true,
|
||||||
transition: TransitionOptions {
|
transition: TransitionOptions {
|
||||||
status: TRANSITION_PENDING.to_string(),
|
status: TRANSITION_PENDING.to_string(),
|
||||||
tier: tier_name.clone(),
|
tier: tier_name.clone(),
|
||||||
etag: original.etag.clone().unwrap_or_default(),
|
etag: original.etag.clone().unwrap_or_default(),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
version_id: original.version_id.map(|version| version.to_string()),
|
|
||||||
mod_time: original.mod_time,
|
mod_time: original.mod_time,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
@@ -7014,7 +7236,10 @@ mod transition_source_identity_matrix_tests {
|
|||||||
|
|
||||||
let mut changed = source.clone();
|
let mut changed = source.clone();
|
||||||
match field {
|
match field {
|
||||||
IdentityField::VersionId => changed.version_id = Some(Uuid::new_v4()),
|
IdentityField::VersionId => {
|
||||||
|
changed.version_id = Some(Uuid::new_v4());
|
||||||
|
changed.fresh = true;
|
||||||
|
}
|
||||||
IdentityField::DataDir => changed.data_dir = Some(Uuid::new_v4()),
|
IdentityField::DataDir => changed.data_dir = Some(Uuid::new_v4()),
|
||||||
IdentityField::ModTime => {
|
IdentityField::ModTime => {
|
||||||
changed.mod_time = changed.mod_time.map(|value| value + time::Duration::nanoseconds(1));
|
changed.mod_time = changed.mod_time.map(|value| value + time::Duration::nanoseconds(1));
|
||||||
@@ -7032,12 +7257,19 @@ mod transition_source_identity_matrix_tests {
|
|||||||
.await
|
.await
|
||||||
.expect("single-field metadata drift should be written");
|
.expect("single-field metadata drift should be written");
|
||||||
}
|
}
|
||||||
|
let persisted_opts = ObjectOptions {
|
||||||
|
version_id: changed.version_id.map(|version_id| version_id.to_string()),
|
||||||
|
versioned: true,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let (persisted, _, _) = set_disks
|
||||||
|
.get_object_fileinfo(bucket, &object, &persisted_opts, true, false)
|
||||||
|
.await
|
||||||
|
.expect("drifted source metadata should resolve");
|
||||||
put_barrier.release();
|
put_barrier.release();
|
||||||
|
|
||||||
transition
|
let result = transition.await.expect("transition task should not panic");
|
||||||
.await
|
assert!(result.is_err(), "transition must reject {field:?} drift");
|
||||||
.expect("transition task should not panic")
|
|
||||||
.expect_err("transition must reject a source whose identity changed after upload");
|
|
||||||
let expected_attempts = index + 1;
|
let expected_attempts = index + 1;
|
||||||
assert_eq!(backend.put_count().await, expected_attempts);
|
assert_eq!(backend.put_count().await, expected_attempts);
|
||||||
assert_eq!(backend.remove_count().await, expected_attempts);
|
assert_eq!(backend.remove_count().await, expected_attempts);
|
||||||
@@ -7048,26 +7280,26 @@ mod transition_source_identity_matrix_tests {
|
|||||||
);
|
);
|
||||||
|
|
||||||
match field {
|
match field {
|
||||||
IdentityField::VersionId => assert_ne!(source.version_id, changed.version_id),
|
IdentityField::VersionId => assert_ne!(source.version_id, persisted.version_id),
|
||||||
IdentityField::DataDir => assert_ne!(source.data_dir, changed.data_dir),
|
IdentityField::DataDir => assert_ne!(source.data_dir, persisted.data_dir),
|
||||||
IdentityField::ModTime => assert_ne!(source.mod_time, changed.mod_time),
|
IdentityField::ModTime => assert_ne!(source.mod_time, persisted.mod_time),
|
||||||
IdentityField::Size => assert_ne!(source.size, changed.size),
|
IdentityField::Size => assert_ne!(source.size, persisted.size),
|
||||||
IdentityField::Etag => assert_ne!(get_raw_etag(&source.metadata), get_raw_etag(&changed.metadata)),
|
IdentityField::Etag => assert_ne!(get_raw_etag(&source.metadata), get_raw_etag(&persisted.metadata)),
|
||||||
}
|
}
|
||||||
if !matches!(field, IdentityField::VersionId) {
|
if !matches!(field, IdentityField::VersionId) {
|
||||||
assert_eq!(source.version_id, changed.version_id);
|
assert_eq!(source.version_id, persisted.version_id);
|
||||||
}
|
}
|
||||||
if !matches!(field, IdentityField::DataDir) {
|
if !matches!(field, IdentityField::DataDir) {
|
||||||
assert_eq!(source.data_dir, changed.data_dir);
|
assert_eq!(source.data_dir, persisted.data_dir);
|
||||||
}
|
}
|
||||||
if !matches!(field, IdentityField::ModTime) {
|
if !matches!(field, IdentityField::ModTime) {
|
||||||
assert_eq!(source.mod_time, changed.mod_time);
|
assert_eq!(source.mod_time, persisted.mod_time);
|
||||||
}
|
}
|
||||||
if !matches!(field, IdentityField::Size) {
|
if !matches!(field, IdentityField::Size) {
|
||||||
assert_eq!(source.size, changed.size);
|
assert_eq!(source.size, persisted.size);
|
||||||
}
|
}
|
||||||
if !matches!(field, IdentityField::Etag) {
|
if !matches!(field, IdentityField::Etag) {
|
||||||
assert_eq!(get_raw_etag(&source.metadata), get_raw_etag(&changed.metadata));
|
assert_eq!(get_raw_etag(&source.metadata), get_raw_etag(&persisted.metadata));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,7 +13,10 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
use crate::bucket::lifecycle::lifecycle;
|
||||||
|
use rustfs_filemeta::RestoreStatusOps;
|
||||||
use rustfs_utils::http::headers::{AMZ_RESTORE_EXPIRY_DAYS, AMZ_RESTORE_REQUEST_DATE};
|
use rustfs_utils::http::headers::{AMZ_RESTORE_EXPIRY_DAYS, AMZ_RESTORE_REQUEST_DATE};
|
||||||
|
use s3s::dto::{RestoreStatus, Timestamp};
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
struct RestoreCleanupIdentity {
|
struct RestoreCleanupIdentity {
|
||||||
@@ -43,6 +46,69 @@ impl RestoreCleanupIdentity {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl SetDisks {
|
impl SetDisks {
|
||||||
|
pub(super) async fn finalize_restore_metadata(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
obj_info: &ObjectInfo,
|
||||||
|
opts: &ObjectOptions,
|
||||||
|
) -> Result<ObjectInfo> {
|
||||||
|
let expected = RestoreCleanupIdentity::from_object_info(obj_info);
|
||||||
|
let expected_operation_id = restore_operation_id_from_metadata(&opts.user_defined)?;
|
||||||
|
let expected_etag = obj_info
|
||||||
|
.etag
|
||||||
|
.clone()
|
||||||
|
.unwrap_or_else(|| get_raw_etag(obj_info.user_defined.as_ref()));
|
||||||
|
let version_id = expected.version_id.map(|v| v.to_string());
|
||||||
|
let _lock_guard = if !opts.no_lock {
|
||||||
|
Some(
|
||||||
|
self.acquire_write_lock_diag("restore_finalize_metadata", bucket, object)
|
||||||
|
.await?,
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
};
|
||||||
|
let read_opts = ObjectOptions {
|
||||||
|
version_id,
|
||||||
|
versioned: opts.versioned,
|
||||||
|
version_suspended: opts.version_suspended,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let (mut fi, _, disks) = self
|
||||||
|
.get_object_fileinfo_gated(bucket, object, &read_opts, false, false)
|
||||||
|
.await?;
|
||||||
|
if let Some(expected_operation_id) = expected_operation_id {
|
||||||
|
require_restore_operation_id(&fi.metadata, expected_operation_id)?;
|
||||||
|
}
|
||||||
|
if !expected.matches_file_info(&fi, &expected_etag) {
|
||||||
|
return Err(Error::other("restored object changed before restore metadata finalization"));
|
||||||
|
}
|
||||||
|
let restore_expiry =
|
||||||
|
lifecycle::expected_expiry_time(OffsetDateTime::now_utc(), opts.transition.restore_request.days.unwrap_or(1));
|
||||||
|
fi.metadata.insert(
|
||||||
|
X_AMZ_RESTORE.as_str().to_string(),
|
||||||
|
RestoreStatus {
|
||||||
|
is_restore_in_progress: Some(false),
|
||||||
|
restore_expiry_date: Some(Timestamp::from(restore_expiry)),
|
||||||
|
}
|
||||||
|
.to_string(),
|
||||||
|
);
|
||||||
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
|
self.update_object_meta_with_opts(
|
||||||
|
bucket,
|
||||||
|
object,
|
||||||
|
fi.clone(),
|
||||||
|
disks.as_slice(),
|
||||||
|
&UpdateMetadataOpts {
|
||||||
|
replace_user_metadata: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||||
|
Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended))
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn update_restore_metadata(
|
pub async fn update_restore_metadata(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
|
|||||||
@@ -13,10 +13,11 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::bucket::lifecycle::lifecycle::{TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions};
|
use crate::bucket::lifecycle::lifecycle::{TRANSITION_COMPLETE, TRANSITION_PENDING, TransitionOptions, expected_expiry_time};
|
||||||
use crate::ecstore_validation_blackbox::make_local_set_disks;
|
use crate::ecstore_validation_blackbox::make_local_set_disks;
|
||||||
use crate::services::tier::test_util::register_mock_tier;
|
use crate::services::tier::test_util::register_mock_tier;
|
||||||
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
|
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
|
||||||
|
use rustfs_filemeta::{RestoreStatusOps as _, parse_restore_obj_status};
|
||||||
use tokio::io::AsyncReadExt;
|
use tokio::io::AsyncReadExt;
|
||||||
|
|
||||||
async fn prime_metadata_generation(set_disks: &SetDisks, bucket: &str, object: &str) -> GetObjectMetadataCacheKey {
|
async fn prime_metadata_generation(set_disks: &SetDisks, bucket: &str, object: &str) -> GetObjectMetadataCacheKey {
|
||||||
@@ -83,13 +84,57 @@ async fn transition_and_restore_reclaim_prior_metadata_generations() {
|
|||||||
let transitioned_generation = prime_metadata_generation(&set_disks, bucket, object).await;
|
let transitioned_generation = prime_metadata_generation(&set_disks, bucket, object).await;
|
||||||
let mut restore_opts = ObjectOptions::default();
|
let mut restore_opts = ObjectOptions::default();
|
||||||
restore_opts.transition.restore_request.days = Some(1);
|
restore_opts.transition.restore_request.days = Some(1);
|
||||||
Arc::clone(&set_disks)
|
let restore_started = OffsetDateTime::now_utc();
|
||||||
.restore_transitioned_object(bucket, object, &restore_opts)
|
let expiry_from_restore_start = temp_env::async_with_vars(
|
||||||
.await
|
[
|
||||||
.expect("restore should succeed");
|
("RUSTFS_ILM_DEBUG_DAY_SECS", Some("1")),
|
||||||
|
("RUSTFS_ILM_PROCESS_TIME", Some("1")),
|
||||||
|
],
|
||||||
|
async {
|
||||||
|
let expiry_from_restore_start = expected_expiry_time(restore_started, 1);
|
||||||
|
let get_barrier = backend.arm_get_barrier().await;
|
||||||
|
let restore_set = Arc::clone(&set_disks);
|
||||||
|
let restore =
|
||||||
|
tokio::spawn(async move { restore_set.restore_transitioned_object(bucket, object, &restore_opts).await });
|
||||||
|
get_barrier.wait_until_paused().await;
|
||||||
|
tokio::time::timeout(Duration::from_secs(5), async {
|
||||||
|
loop {
|
||||||
|
if expected_expiry_time(OffsetDateTime::now_utc(), 1) > expiry_from_restore_start {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("test clock should cross the next accelerated lifecycle boundary");
|
||||||
|
get_barrier.release();
|
||||||
|
restore
|
||||||
|
.await
|
||||||
|
.expect("restore task should join")
|
||||||
|
.expect("restore should succeed");
|
||||||
|
expiry_from_restore_start
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await;
|
||||||
assert_generation_reclaimed(&set_disks, &transitioned_generation).await;
|
assert_generation_reclaimed(&set_disks, &transitioned_generation).await;
|
||||||
assert_eq!(backend.get_count().await, 1, "restore should read the remote candidate exactly once");
|
assert_eq!(backend.get_count().await, 1, "restore should read the remote candidate exactly once");
|
||||||
|
|
||||||
|
let restored_info = set_disks
|
||||||
|
.get_object_info(bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("restored object metadata should be readable");
|
||||||
|
let restore_status = parse_restore_obj_status(
|
||||||
|
restored_info
|
||||||
|
.user_defined
|
||||||
|
.get(s3s::header::X_AMZ_RESTORE.as_str())
|
||||||
|
.expect("completed restore header should be present"),
|
||||||
|
)
|
||||||
|
.expect("completed restore header should parse");
|
||||||
|
assert!(
|
||||||
|
restore_status.expiry().expect("completed restore should have an expiry") > expiry_from_restore_start,
|
||||||
|
"restore expiry must be based on completion, not the time the remote copy started"
|
||||||
|
);
|
||||||
|
|
||||||
let mut restored = Vec::new();
|
let mut restored = Vec::new();
|
||||||
set_disks
|
set_disks
|
||||||
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||||
|
|||||||
@@ -1311,14 +1311,16 @@ mod tests {
|
|||||||
version_id: "version-a".to_string(),
|
version_id: "version-a".to_string(),
|
||||||
tier_name: tier_a.to_string(),
|
tier_name: tier_a.to_string(),
|
||||||
backend_identity: Some(identity_a),
|
backend_identity: Some(identity_a),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
};
|
};
|
||||||
let entry_b = Jentry {
|
let entry_b = Jentry {
|
||||||
obj_name: "remote-b".to_string(),
|
obj_name: "remote-b".to_string(),
|
||||||
version_id: "version-b".to_string(),
|
version_id: "version-b".to_string(),
|
||||||
tier_name: tier_b.to_string(),
|
tier_name: tier_b.to_string(),
|
||||||
backend_identity: Some(identity_b),
|
backend_identity: Some(identity_b),
|
||||||
version_id_exact: false,
|
version_id_exact: true,
|
||||||
|
version_state: rustfs_filemeta::TransitionVersionState::Exact,
|
||||||
};
|
};
|
||||||
let remove_a = backend_a.arm_failing_remove_barrier().await;
|
let remove_a = backend_a.arm_failing_remove_barrier().await;
|
||||||
persist_tier_delete_journal_entry(store_a.clone(), &entry_a)
|
persist_tier_delete_journal_entry(store_a.clone(), &entry_a)
|
||||||
@@ -2720,10 +2722,12 @@ mod tests {
|
|||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
async fn transition_transaction_recovery_deletes_provider_recovered_unknown_upload() {
|
async fn transition_transaction_recovery_deletes_provider_recovered_unknown_upload() {
|
||||||
let versioned_remote = uuid::Uuid::new_v4().to_string();
|
let versioned_remote = uuid::Uuid::new_v4().to_string();
|
||||||
|
let nil_remote = uuid::Uuid::nil().to_string();
|
||||||
for (case, tier_name, remote_version) in [
|
for (case, tier_name, remote_version) in [
|
||||||
("missing", "TXPROBEMISSING", None),
|
("missing", "TXPROBEMISSING", None),
|
||||||
("unversioned", "TXPROBEUNVERSIONED", Some(String::new())),
|
("unversioned", "TXPROBEUNVERSIONED", Some(String::new())),
|
||||||
("versioned", "TXPROBEVERSIONED", Some(versioned_remote)),
|
("versioned", "TXPROBEVERSIONED", Some(versioned_remote)),
|
||||||
|
("nil-version", "TXPROBENILVERSION", Some(nil_remote)),
|
||||||
] {
|
] {
|
||||||
let temp_dir = tempfile::tempdir().expect("create temp store dir");
|
let temp_dir = tempfile::tempdir().expect("create temp store dir");
|
||||||
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store(
|
let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||||
|
|||||||
@@ -219,6 +219,16 @@ impl ErasureInfo {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// #[derive(Debug, Clone)]
|
// #[derive(Debug, Clone)]
|
||||||
|
#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Clone, Copy, Default)]
|
||||||
|
#[serde(rename_all = "kebab-case")]
|
||||||
|
pub enum TransitionVersionState {
|
||||||
|
#[default]
|
||||||
|
Unknown,
|
||||||
|
KnownDisabled,
|
||||||
|
SuspendedNull,
|
||||||
|
Exact,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)]
|
#[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)]
|
||||||
pub struct FileInfo {
|
pub struct FileInfo {
|
||||||
pub volume: String,
|
pub volume: String,
|
||||||
@@ -230,6 +240,10 @@ pub struct FileInfo {
|
|||||||
pub transitioned_objname: String,
|
pub transitioned_objname: String,
|
||||||
pub transition_tier: String,
|
pub transition_tier: String,
|
||||||
pub transition_version_id: Option<Uuid>,
|
pub transition_version_id: Option<Uuid>,
|
||||||
|
#[serde(default)]
|
||||||
|
pub transition_version: Option<String>,
|
||||||
|
#[serde(default)]
|
||||||
|
pub transition_version_state: TransitionVersionState,
|
||||||
pub expire_restored: bool,
|
pub expire_restored: bool,
|
||||||
pub data_dir: Option<Uuid>,
|
pub data_dir: Option<Uuid>,
|
||||||
pub mod_time: Option<OffsetDateTime>,
|
pub mod_time: Option<OffsetDateTime>,
|
||||||
@@ -459,6 +473,10 @@ impl FileInfo {
|
|||||||
if self.mod_time.is_none_or(|mod_time| mod_time <= OffsetDateTime::UNIX_EPOCH)
|
if self.mod_time.is_none_or(|mod_time| mod_time <= OffsetDateTime::UNIX_EPOCH)
|
||||||
|| (!allow_nil_version_id && self.version_id.is_some_and(|version_id| version_id.is_nil()))
|
|| (!allow_nil_version_id && self.version_id.is_some_and(|version_id| version_id.is_nil()))
|
||||||
|| self.transition_version_id.is_some_and(|version_id| version_id.is_nil())
|
|| self.transition_version_id.is_some_and(|version_id| version_id.is_nil())
|
||||||
|
|| self
|
||||||
|
.transition_version
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|version_id| version_id.is_empty())
|
||||||
|| self.size != 0
|
|| self.size != 0
|
||||||
|| self.data_dir.is_some()
|
|| self.data_dir.is_some()
|
||||||
|| self.mode.is_some()
|
|| self.mode.is_some()
|
||||||
@@ -492,6 +510,7 @@ impl FileInfo {
|
|||||||
|| !self.transitioned_objname.is_empty()
|
|| !self.transitioned_objname.is_empty()
|
||||||
|| !self.transition_tier.is_empty()
|
|| !self.transition_tier.is_empty()
|
||||||
|| self.transition_version_id.is_some()
|
|| self.transition_version_id.is_some()
|
||||||
|
|| self.transition_version.is_some()
|
||||||
|| self.expire_restored
|
|| self.expire_restored
|
||||||
|| self.size != 0
|
|| self.size != 0
|
||||||
|| self.data_dir.is_some()
|
|| self.data_dir.is_some()
|
||||||
@@ -536,6 +555,25 @@ impl FileInfo {
|
|||||||
/// return `None`.
|
/// return `None`.
|
||||||
pub fn validate(&self, mode: ValidationMode) -> Result<Option<ValidatedErasureLayout>> {
|
pub fn validate(&self, mode: ValidationMode) -> Result<Option<ValidatedErasureLayout>> {
|
||||||
self.validate_collection_bounds()?;
|
self.validate_collection_bounds()?;
|
||||||
|
if let (Some(version), Some(version_id)) = (&self.transition_version, self.transition_version_id)
|
||||||
|
&& Uuid::parse_str(version).ok() != Some(version_id)
|
||||||
|
{
|
||||||
|
return Err(Error::FileCorrupt);
|
||||||
|
}
|
||||||
|
let transition_state_valid = match self.transition_version_state {
|
||||||
|
TransitionVersionState::Unknown => true,
|
||||||
|
TransitionVersionState::KnownDisabled => self.transition_version.is_none() && self.transition_version_id.is_none(),
|
||||||
|
TransitionVersionState::SuspendedNull => {
|
||||||
|
self.transition_version.as_deref() == Some("null") && self.transition_version_id.is_none()
|
||||||
|
}
|
||||||
|
TransitionVersionState::Exact => self
|
||||||
|
.transition_version
|
||||||
|
.as_deref()
|
||||||
|
.is_some_and(|version| version != "null" && !version.is_empty()),
|
||||||
|
};
|
||||||
|
if !transition_state_valid {
|
||||||
|
return Err(Error::FileCorrupt);
|
||||||
|
}
|
||||||
|
|
||||||
let erasure_layout = match mode {
|
let erasure_layout = match mode {
|
||||||
ValidationMode::RequireErasure => Some(self.validate_erasure_geometry()?),
|
ValidationMode::RequireErasure => Some(self.validate_erasure_geometry()?),
|
||||||
@@ -832,6 +870,8 @@ impl FileInfo {
|
|||||||
&& self.transition_tier == other.transition_tier
|
&& self.transition_tier == other.transition_tier
|
||||||
&& self.transitioned_objname == other.transitioned_objname
|
&& self.transitioned_objname == other.transitioned_objname
|
||||||
&& self.transition_version_id == other.transition_version_id
|
&& self.transition_version_id == other.transition_version_id
|
||||||
|
&& self.transition_version == other.transition_version
|
||||||
|
&& self.transition_version_state == other.transition_version_state
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Check if metadata maps are equal
|
/// Check if metadata maps are equal
|
||||||
@@ -1351,6 +1391,15 @@ mod tests {
|
|||||||
assert_file_corrupt(&fi, ValidationMode::DeleteOnly);
|
assert_file_corrupt(&fi, ValidationMode::DeleteOnly);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn metadata_read_validation_rejects_conflicting_transition_versions() {
|
||||||
|
let mut fi = one_shard_validation_fileinfo(1);
|
||||||
|
fi.transition_version_id = Some(Uuid::new_v4());
|
||||||
|
fi.transition_version = Some(Uuid::new_v4().to_string());
|
||||||
|
|
||||||
|
assert_file_corrupt(&fi, ValidationMode::RequireErasure);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn metadata_read_validation_requires_canonical_delete_marker_shape() {
|
fn metadata_read_validation_requires_canonical_delete_marker_shape() {
|
||||||
let marker = FileInfo {
|
let marker = FileInfo {
|
||||||
@@ -1722,6 +1771,12 @@ mod tests {
|
|||||||
transitioned_objname,
|
transitioned_objname,
|
||||||
transition_tier,
|
transition_tier,
|
||||||
transition_version_id,
|
transition_version_id,
|
||||||
|
transition_version: transition_version_id.map(|version_id| version_id.to_string()),
|
||||||
|
transition_version_state: if transition_version_id.is_some() {
|
||||||
|
TransitionVersionState::Exact
|
||||||
|
} else {
|
||||||
|
TransitionVersionState::Unknown
|
||||||
|
},
|
||||||
expire_restored,
|
expire_restored,
|
||||||
data_dir,
|
data_dir,
|
||||||
mod_time,
|
mod_time,
|
||||||
|
|||||||
@@ -26,13 +26,14 @@ use super::msgp_decode::{
|
|||||||
PrependByteReader, prealloc_hint, read_exact_vec, read_nil_or_array_len, read_nil_or_map_len, skip_msgp_value,
|
PrependByteReader, prealloc_hint, read_exact_vec, read_nil_or_array_len, read_nil_or_map_len, skip_msgp_value,
|
||||||
};
|
};
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::ChecksumInfo;
|
use crate::{ChecksumInfo, TransitionVersionState};
|
||||||
use rustfs_utils::HashAlgorithm;
|
use rustfs_utils::HashAlgorithm;
|
||||||
use rustfs_utils::http::{
|
use rustfs_utils::http::{
|
||||||
RUSTFS_INTERNAL_PREFIX, SUFFIX_CRC, SUFFIX_FREE_VERSION, SUFFIX_INLINE_DATA, SUFFIX_PURGESTATUS, SUFFIX_TIER_FV_ID,
|
RUSTFS_INTERNAL_PREFIX, SUFFIX_CRC, SUFFIX_FREE_VERSION, SUFFIX_INLINE_DATA, SUFFIX_PURGESTATUS, SUFFIX_TIER_FV_ID,
|
||||||
SUFFIX_TIER_FV_MARKER, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
SUFFIX_TIER_FV_MARKER, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||||
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, contains_key_bytes, get_bytes, get_consistent_bytes, get_str,
|
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_bytes,
|
||||||
has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, strip_internal_prefix,
|
get_bytes, get_consistent_bytes, get_str, has_internal_suffix, insert_bytes, is_internal_key, remove_bytes,
|
||||||
|
strip_internal_prefix,
|
||||||
};
|
};
|
||||||
|
|
||||||
const MSGPACK_EXT8: u8 = 0xc7;
|
const MSGPACK_EXT8: u8 = 0xc7;
|
||||||
@@ -43,6 +44,7 @@ const MSGPACK_FIXEXT8: u8 = 0xd7;
|
|||||||
const MSGPACK_TIME_EXT_LEGACY: i8 = 5;
|
const MSGPACK_TIME_EXT_LEGACY: i8 = 5;
|
||||||
const MSGPACK_TIME_EXT_OFFICIAL: i8 = -1;
|
const MSGPACK_TIME_EXT_OFFICIAL: i8 = -1;
|
||||||
const MSGPACK_TIME_LEN: u8 = 12;
|
const MSGPACK_TIME_LEN: u8 = 12;
|
||||||
|
const MAX_TRANSITION_VERSION_LEN: usize = 1024;
|
||||||
|
|
||||||
/// Sentinel signature returned when a version has no computable body (invalid /
|
/// Sentinel signature returned when a version has no computable body (invalid /
|
||||||
/// missing inner object). Mirrors MinIO's `signatureErr` so such versions never
|
/// missing inner object). Mirrors MinIO's `signatureErr` so such versions never
|
||||||
@@ -251,23 +253,93 @@ fn parse_legacy_uuid_bytes(bytes: &[u8], field: &str) -> Result<Option<Uuid>> {
|
|||||||
|
|
||||||
/// Decode a stored transitioned-version-id from a version's `meta_sys`.
|
/// Decode a stored transitioned-version-id from a version's `meta_sys`.
|
||||||
///
|
///
|
||||||
/// RustFS writes it as 16 raw UUID bytes; MinIO-migrated tiered objects store
|
/// Legacy RustFS writes used 16 raw UUID bytes. New writes and MinIO-migrated
|
||||||
/// the remote tier's version id as a UUID *string*. Accept both, and treat any
|
/// records use the provider's exact UTF-8 version text. Empty, nil UUID, and
|
||||||
/// absent / nil / otherwise-unparseable value as "no tier version" (matching the
|
/// malformed bytes are not usable remote versions.
|
||||||
/// tolerant pre-hardening behavior) rather than failing the whole object read —
|
fn transitioned_version_from_meta_sys(meta_sys: &HashMap<String, Vec<u8>>) -> Result<Option<String>> {
|
||||||
/// a malformed tier id must not make an otherwise-readable object unreadable.
|
if !contains_key_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID) {
|
||||||
fn transitioned_version_id_from_meta_sys(meta_sys: &HashMap<String, Vec<u8>>) -> Option<Uuid> {
|
return Ok(None);
|
||||||
let value = get_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)?;
|
}
|
||||||
|
let Some(value) = get_consistent_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID) else {
|
||||||
|
return Ok(None);
|
||||||
|
};
|
||||||
|
let value = value.to_vec();
|
||||||
if value.is_empty() {
|
if value.is_empty() {
|
||||||
return None;
|
return Ok(None);
|
||||||
}
|
}
|
||||||
if let Ok(id) = Uuid::from_slice(&value) {
|
if let Ok(id) = Uuid::from_slice(&value) {
|
||||||
return (!id.is_nil()).then_some(id);
|
return Ok((!id.is_nil()).then(|| id.to_string()));
|
||||||
}
|
}
|
||||||
std::str::from_utf8(&value)
|
let Ok(value) = String::from_utf8(value) else {
|
||||||
|
return Ok(None);
|
||||||
|
};
|
||||||
|
if value.is_empty()
|
||||||
|
|| value.len() > MAX_TRANSITION_VERSION_LEN
|
||||||
|
|| value.chars().any(char::is_control)
|
||||||
|
|| Uuid::parse_str(&value).is_ok_and(|id| id.is_nil())
|
||||||
|
{
|
||||||
|
Ok(None)
|
||||||
|
} else {
|
||||||
|
Ok(Some(value))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn transition_version_state_from_meta_sys(
|
||||||
|
meta_sys: &HashMap<String, Vec<u8>>,
|
||||||
|
version: Option<&str>,
|
||||||
|
) -> Result<TransitionVersionState> {
|
||||||
|
if !contains_key_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE) {
|
||||||
|
return Ok(TransitionVersionState::Unknown);
|
||||||
|
}
|
||||||
|
let value = get_consistent_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE).ok_or(Error::FileCorrupt)?;
|
||||||
|
let state = match value {
|
||||||
|
b"known-disabled" => TransitionVersionState::KnownDisabled,
|
||||||
|
b"suspended-null" => TransitionVersionState::SuspendedNull,
|
||||||
|
b"exact" => TransitionVersionState::Exact,
|
||||||
|
b"unknown" => TransitionVersionState::Unknown,
|
||||||
|
_ => return Err(Error::FileCorrupt),
|
||||||
|
};
|
||||||
|
let valid = match state {
|
||||||
|
TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(),
|
||||||
|
TransitionVersionState::SuspendedNull => version == Some("null"),
|
||||||
|
TransitionVersionState::Exact => version.is_some_and(|value| value != "null"),
|
||||||
|
};
|
||||||
|
valid.then_some(state).ok_or(Error::FileCorrupt)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn transition_version_state_bytes(state: TransitionVersionState) -> &'static [u8] {
|
||||||
|
match state {
|
||||||
|
TransitionVersionState::Unknown => b"unknown",
|
||||||
|
TransitionVersionState::KnownDisabled => b"known-disabled",
|
||||||
|
TransitionVersionState::SuspendedNull => b"suspended-null",
|
||||||
|
TransitionVersionState::Exact => b"exact",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn set_transition_version_state(meta_sys: &mut HashMap<String, Vec<u8>>, state: TransitionVersionState) {
|
||||||
|
if state == TransitionVersionState::Unknown {
|
||||||
|
remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE);
|
||||||
|
} else {
|
||||||
|
insert_bytes(
|
||||||
|
meta_sys,
|
||||||
|
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
|
transition_version_state_bytes(state).to_vec(),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn legacy_transitioned_version_id_from_meta_sys(meta_sys: &HashMap<String, Vec<u8>>) -> Option<Uuid> {
|
||||||
|
transitioned_version_from_meta_sys(meta_sys)
|
||||||
.ok()
|
.ok()
|
||||||
.and_then(|s| Uuid::parse_str(s.trim()).ok())
|
.flatten()
|
||||||
.filter(|id| !id.is_nil())
|
.and_then(|value| Uuid::parse_str(&value).ok())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn transitioned_version_bytes(fi: &FileInfo) -> Option<Vec<u8>> {
|
||||||
|
fi.transition_version
|
||||||
|
.as_ref()
|
||||||
|
.map(|version| version.as_bytes().to_vec())
|
||||||
|
.or_else(|| fi.transition_version_id.map(|version_id| version_id.as_bytes().to_vec()))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn parse_legacy_erasure_algo(value: &str) -> ErasureAlgo {
|
fn parse_legacy_erasure_algo(value: &str) -> ErasureAlgo {
|
||||||
@@ -2398,7 +2470,9 @@ impl MetaObject {
|
|||||||
let transitioned_objname = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_OBJECTNAME)
|
let transitioned_objname = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_OBJECTNAME)
|
||||||
.map(|v| String::from_utf8_lossy(&v).to_string())
|
.map(|v| String::from_utf8_lossy(&v).to_string())
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
let transition_version_id = transitioned_version_id_from_meta_sys(&self.meta_sys);
|
let transition_version = transitioned_version_from_meta_sys(&self.meta_sys)?;
|
||||||
|
let transition_version_state = transition_version_state_from_meta_sys(&self.meta_sys, transition_version.as_deref())?;
|
||||||
|
let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok());
|
||||||
let transition_tier = get_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER)
|
let transition_tier = get_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER)
|
||||||
.map(|v| String::from_utf8_lossy(&v).to_string())
|
.map(|v| String::from_utf8_lossy(&v).to_string())
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
@@ -2419,6 +2493,8 @@ impl MetaObject {
|
|||||||
transition_status,
|
transition_status,
|
||||||
transitioned_objname,
|
transitioned_objname,
|
||||||
transition_version_id,
|
transition_version_id,
|
||||||
|
transition_version,
|
||||||
|
transition_version_state,
|
||||||
transition_tier,
|
transition_tier,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
})
|
})
|
||||||
@@ -2431,13 +2507,12 @@ impl MetaObject {
|
|||||||
SUFFIX_TRANSITIONED_OBJECTNAME,
|
SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||||
fi.transitioned_objname.as_bytes().to_vec(),
|
fi.transitioned_objname.as_bytes().to_vec(),
|
||||||
);
|
);
|
||||||
if let Some(transition_version_id) = fi.transition_version_id.as_ref() {
|
if let Some(transition_version) = transitioned_version_bytes(fi) {
|
||||||
insert_bytes(
|
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
||||||
&mut self.meta_sys,
|
} else {
|
||||||
SUFFIX_TRANSITIONED_VERSION_ID,
|
remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID);
|
||||||
transition_version_id.as_bytes().to_vec(),
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
set_transition_version_state(&mut self.meta_sys, fi.transition_version_state);
|
||||||
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
|
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
|
||||||
if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
|
if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
|
||||||
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes());
|
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes());
|
||||||
@@ -2501,6 +2576,7 @@ impl MetaObject {
|
|||||||
SUFFIX_TRANSITION_TIER,
|
SUFFIX_TRANSITION_TIER,
|
||||||
SUFFIX_TRANSITIONED_OBJECTNAME,
|
SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||||
SUFFIX_TRANSITIONED_VERSION_ID,
|
SUFFIX_TRANSITIONED_VERSION_ID,
|
||||||
|
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||||
] {
|
] {
|
||||||
if let Some(v) = get_bytes(&self.meta_sys, suffix) {
|
if let Some(v) = get_bytes(&self.meta_sys, suffix) {
|
||||||
insert_bytes(&mut delete_marker.meta_sys, suffix, v);
|
insert_bytes(&mut delete_marker.meta_sys, suffix, v);
|
||||||
@@ -2562,8 +2638,11 @@ impl From<FileInfo> for MetaObject {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Some(vid) = &value.transition_version_id {
|
if let Some(transition_version) = transitioned_version_bytes(&value) {
|
||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, vid.as_bytes().to_vec());
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
||||||
|
}
|
||||||
|
if !value.transition_status.is_empty() {
|
||||||
|
set_transition_version_state(&mut meta_sys, value.transition_version_state);
|
||||||
}
|
}
|
||||||
|
|
||||||
if !value.transition_tier.is_empty() {
|
if !value.transition_tier.is_empty() {
|
||||||
@@ -2706,7 +2785,11 @@ impl MetaDeleteMarker {
|
|||||||
.map(|v| String::from_utf8_lossy(&v).to_string())
|
.map(|v| String::from_utf8_lossy(&v).to_string())
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
|
|
||||||
fi.transition_version_id = transitioned_version_id_from_meta_sys(&self.meta_sys);
|
fi.transition_version = transitioned_version_from_meta_sys(&self.meta_sys).ok().flatten();
|
||||||
|
fi.transition_version_id = legacy_transitioned_version_id_from_meta_sys(&self.meta_sys);
|
||||||
|
fi.transition_version_state =
|
||||||
|
transition_version_state_from_meta_sys(&self.meta_sys, fi.transition_version.as_deref())
|
||||||
|
.unwrap_or(TransitionVersionState::Unknown);
|
||||||
}
|
}
|
||||||
|
|
||||||
fi
|
fi
|
||||||
@@ -2859,8 +2942,11 @@ impl From<FileInfo> for MetaDeleteMarker {
|
|||||||
value.transitioned_objname.as_bytes().to_vec(),
|
value.transitioned_objname.as_bytes().to_vec(),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
if let Some(version_id) = value.transition_version_id {
|
if let Some(transition_version) = transitioned_version_bytes(&value) {
|
||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, version_id.as_bytes().to_vec());
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version);
|
||||||
|
}
|
||||||
|
if !value.transition_status.is_empty() || value.tier_free_version() {
|
||||||
|
set_transition_version_state(&mut meta_sys, value.transition_version_state);
|
||||||
}
|
}
|
||||||
if !value.transition_tier.is_empty() {
|
if !value.transition_tier.is_empty() {
|
||||||
insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec());
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec());
|
||||||
@@ -3412,7 +3498,7 @@ mod tests {
|
|||||||
.insert("x-rustfs-internal-healing".to_string(), "true".to_string());
|
.insert("x-rustfs-internal-healing".to_string(), "true".to_string());
|
||||||
marker.metadata.insert("content-type".to_string(), "text/plain".to_string());
|
marker.metadata.insert("content-type".to_string(), "text/plain".to_string());
|
||||||
let remote_version_id = Uuid::new_v4();
|
let remote_version_id = Uuid::new_v4();
|
||||||
marker.transition_version_id = Some(remote_version_id);
|
marker.transition_version = Some(remote_version_id.to_string());
|
||||||
|
|
||||||
let converted = MetaDeleteMarker::from(marker);
|
let converted = MetaDeleteMarker::from(marker);
|
||||||
|
|
||||||
@@ -3420,7 +3506,19 @@ mod tests {
|
|||||||
assert_eq!(converted.meta_sys.get("x-minio-internal-purgestatus"), Some(&b"pending".to_vec()));
|
assert_eq!(converted.meta_sys.get("x-minio-internal-purgestatus"), Some(&b"pending".to_vec()));
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
get_bytes(&converted.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID),
|
get_bytes(&converted.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID),
|
||||||
Some(remote_version_id.as_bytes().to_vec())
|
Some(remote_version_id.to_string().into_bytes())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
converted
|
||||||
|
.meta_sys
|
||||||
|
.get(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}")),
|
||||||
|
Some(&remote_version_id.to_string().into_bytes())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
converted
|
||||||
|
.meta_sys
|
||||||
|
.get(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX)),
|
||||||
|
Some(&remote_version_id.to_string().into_bytes())
|
||||||
);
|
);
|
||||||
assert!(!converted.meta_sys.contains_key("x-rustfs-internal-healing"));
|
assert!(!converted.meta_sys.contains_key("x-rustfs-internal-healing"));
|
||||||
assert!(!converted.meta_sys.contains_key("content-type"));
|
assert!(!converted.meta_sys.contains_key("content-type"));
|
||||||
@@ -4097,19 +4195,129 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("into_fileinfo");
|
.expect("into_fileinfo");
|
||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
|
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn meta_object_transition_version_id_unparseable_stays_readable_as_none() {
|
fn meta_object_transition_version_id_opaque_text_is_preserved() {
|
||||||
// A non-UUID / non-16-byte tier version id must NOT make the object
|
|
||||||
// unreadable; it is tolerated as "no tier version" (compat with
|
|
||||||
// pre-hardening behavior and foreign/edge metadata).
|
|
||||||
let mut sys = HashMap::new();
|
let mut sys = HashMap::new();
|
||||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec());
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec());
|
||||||
let fi = make_meta_object_with_sys(sys)
|
let fi = make_meta_object_with_sys(sys)
|
||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("unparseable transition version id must not fail the object read");
|
.expect("opaque transition version id must decode");
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42"));
|
||||||
|
assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_state_exact_round_trips_dual_keys() {
|
||||||
|
let id = sample_version_id();
|
||||||
|
let expected_version = id.to_string();
|
||||||
|
let fi = FileInfo {
|
||||||
|
transition_status: "complete".to_string(),
|
||||||
|
transition_version: Some(expected_version.clone()),
|
||||||
|
transition_version_state: TransitionVersionState::Exact,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let object = MetaObject::from(fi);
|
||||||
|
assert_eq!(
|
||||||
|
object
|
||||||
|
.meta_sys
|
||||||
|
.get(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"))
|
||||||
|
.map(Vec::as_slice),
|
||||||
|
Some(b"exact".as_slice())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
object
|
||||||
|
.meta_sys
|
||||||
|
.get(&format!(
|
||||||
|
"{}{SUFFIX_TRANSITIONED_VERSION_STATE}",
|
||||||
|
rustfs_utils::http::MINIO_INTERNAL_PREFIX
|
||||||
|
))
|
||||||
|
.map(Vec::as_slice),
|
||||||
|
Some(b"exact".as_slice())
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
legacy_transitioned_version_id_from_meta_sys(&object.meta_sys),
|
||||||
|
Some(id),
|
||||||
|
"UUID exact writes must remain readable by the legacy UUID consumer"
|
||||||
|
);
|
||||||
|
let decoded = object.into_fileinfo("b", "k", false).expect("exact state should round trip");
|
||||||
|
assert_eq!(decoded.transition_version_state, TransitionVersionState::Exact);
|
||||||
|
assert_eq!(decoded.transition_version.as_deref(), Some(expected_version.as_str()));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn set_transition_known_disabled_removes_stale_version_dual_keys() {
|
||||||
|
let mut meta_sys = HashMap::new();
|
||||||
|
insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, b"stale-legacy-version".to_vec());
|
||||||
|
let mut object = make_meta_object_with_sys(meta_sys);
|
||||||
|
object.set_transition(&FileInfo {
|
||||||
|
transition_status: TRANSITION_COMPLETE.to_string(),
|
||||||
|
transitioned_objname: "remote/object".to_string(),
|
||||||
|
transition_version_state: TransitionVersionState::KnownDisabled,
|
||||||
|
transition_tier: "WARM".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
|
||||||
|
assert_eq!(get_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID), None);
|
||||||
|
assert!(
|
||||||
|
!object
|
||||||
|
.meta_sys
|
||||||
|
.contains_key(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}"))
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!object
|
||||||
|
.meta_sys
|
||||||
|
.contains_key(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX))
|
||||||
|
);
|
||||||
|
let decoded = object
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect("known-disabled transition must remain readable after replacing stale metadata");
|
||||||
|
assert_eq!(decoded.transition_version, None);
|
||||||
|
assert_eq!(decoded.transition_version_state, TransitionVersionState::KnownDisabled);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_state_conflict_fails_closed() {
|
||||||
|
let mut sys = HashMap::new();
|
||||||
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, sample_version_id().as_bytes().to_vec());
|
||||||
|
sys.insert(format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"), b"exact".to_vec());
|
||||||
|
sys.insert(
|
||||||
|
format!("{}{SUFFIX_TRANSITIONED_VERSION_STATE}", rustfs_utils::http::MINIO_INTERNAL_PREFIX),
|
||||||
|
b"known-disabled".to_vec(),
|
||||||
|
);
|
||||||
|
|
||||||
|
make_meta_object_with_sys(sys)
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect_err("conflicting state keys must fail closed");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_id_invalid_utf8_yields_none() {
|
||||||
|
let mut sys = HashMap::new();
|
||||||
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, vec![0xff]);
|
||||||
|
let fi = make_meta_object_with_sys(sys)
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect("invalid transition version bytes must not fail the object read");
|
||||||
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(fi.transition_version, None);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn meta_object_transition_version_id_unsafe_text_yields_none() {
|
||||||
|
for value in [b"opaque\0version".to_vec(), vec![b'x'; MAX_TRANSITION_VERSION_LEN + 1]] {
|
||||||
|
let mut sys = HashMap::new();
|
||||||
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, value);
|
||||||
|
let fi = make_meta_object_with_sys(sys)
|
||||||
|
.into_fileinfo("b", "k", false)
|
||||||
|
.expect("unsafe transition version text must not fail the object read");
|
||||||
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(fi.transition_version, None);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4123,6 +4331,7 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false)
|
.into_fileinfo("b", "k", false)
|
||||||
.expect("string-form transition version id must decode");
|
.expect("string-form transition version id must decode");
|
||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4152,16 +4361,14 @@ mod tests {
|
|||||||
}
|
}
|
||||||
.into_fileinfo("b", "k", false);
|
.into_fileinfo("b", "k", false);
|
||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn delete_marker_free_version_transition_version_id_unparseable_stays_readable() {
|
fn delete_marker_free_version_transition_version_id_opaque_text_is_preserved() {
|
||||||
// A malformed tier version id must not make a free-version record corrupt:
|
|
||||||
// it decodes to None and stays readable. Otherwise free-version expiry
|
|
||||||
// fails and the remote-tier object leaks.
|
|
||||||
let mut sys = HashMap::new();
|
let mut sys = HashMap::new();
|
||||||
insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]);
|
insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]);
|
||||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec());
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec());
|
||||||
insert_bytes(&mut sys, SUFFIX_TRANSITION_TIER, b"WARM".to_vec());
|
insert_bytes(&mut sys, SUFFIX_TRANSITION_TIER, b"WARM".to_vec());
|
||||||
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_OBJECTNAME, b"remote-object".to_vec());
|
insert_bytes(&mut sys, SUFFIX_TRANSITIONED_OBJECTNAME, b"remote-object".to_vec());
|
||||||
let fi = MetaDeleteMarker {
|
let fi = MetaDeleteMarker {
|
||||||
@@ -4172,8 +4379,9 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false);
|
.into_fileinfo("b", "k", false);
|
||||||
|
|
||||||
assert_eq!(fi.transition_version_id, None);
|
assert_eq!(fi.transition_version_id, None);
|
||||||
|
assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42"));
|
||||||
fi.validate_for_metadata_read()
|
fi.validate_for_metadata_read()
|
||||||
.expect("free-version record with an unparseable tier id must remain readable");
|
.expect("free-version record with an opaque tier id must remain readable");
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -4193,6 +4401,7 @@ mod tests {
|
|||||||
.into_fileinfo("b", "k", false);
|
.into_fileinfo("b", "k", false);
|
||||||
|
|
||||||
assert_eq!(fi.transition_version_id, Some(id));
|
assert_eq!(fi.transition_version_id, Some(id));
|
||||||
|
assert_eq!(fi.transition_version, Some(id.to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -698,7 +698,7 @@ impl DataUsageCache {
|
|||||||
let mut visited = HashSet::new();
|
let mut visited = HashSet::new();
|
||||||
visited.insert(hash_path(path).key());
|
visited.insert(hash_path(path).key());
|
||||||
let mut flat = self.flatten_with_guard(root, &mut visited, 0);
|
let mut flat = self.flatten_with_guard(root, &mut visited, 0);
|
||||||
if flat.replication_stats.as_ref().is_some_and(|stats| stats.empty()) {
|
if flat.replication_stats.as_ref().is_some_and(|stats| stats.is_empty()) {
|
||||||
flat.replication_stats = None;
|
flat.replication_stats = None;
|
||||||
}
|
}
|
||||||
Some(flat)
|
Some(flat)
|
||||||
@@ -1574,6 +1574,7 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO};
|
use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO};
|
||||||
use crate::{ScannerGetObjectReader, ScannerPutObjReader};
|
use crate::{ScannerGetObjectReader, ScannerPutObjReader};
|
||||||
|
use rustfs_data_usage::{ReplicationAllStats, ReplicationStats};
|
||||||
use serde_json::Value;
|
use serde_json::Value;
|
||||||
use std::io::Cursor;
|
use std::io::Cursor;
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
@@ -2767,6 +2768,55 @@ mod tests {
|
|||||||
assert!(flat.children.is_empty());
|
assert!(flat.children.is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn size_recursive_prunes_empty_and_preserves_threshold_replication_stats() {
|
||||||
|
let root = hash_path("bucket");
|
||||||
|
let child = hash_path("bucket/child");
|
||||||
|
let mut cache = DataUsageCache::default();
|
||||||
|
cache.replace_hashed(&root, &None, &DataUsageEntry::default());
|
||||||
|
cache.replace_hashed(
|
||||||
|
&child,
|
||||||
|
&Some(root.clone()),
|
||||||
|
&DataUsageEntry {
|
||||||
|
replication_stats: Some(ReplicationAllStats::default()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
cache
|
||||||
|
.size_recursive("bucket")
|
||||||
|
.expect("scanner bucket usage should flatten")
|
||||||
|
.replication_stats
|
||||||
|
.is_none()
|
||||||
|
);
|
||||||
|
|
||||||
|
cache.replace_hashed(
|
||||||
|
&child,
|
||||||
|
&Some(root.clone()),
|
||||||
|
&DataUsageEntry {
|
||||||
|
replication_stats: Some(ReplicationAllStats {
|
||||||
|
targets: HashMap::from([(
|
||||||
|
"arn:test:threshold".to_string(),
|
||||||
|
ReplicationStats {
|
||||||
|
after_threshold_count: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)]),
|
||||||
|
..Default::default()
|
||||||
|
}),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let flattened = cache.size_recursive("bucket").expect("scanner bucket usage should flatten");
|
||||||
|
let replication = flattened
|
||||||
|
.replication_stats
|
||||||
|
.expect("threshold-only replication stats must survive pruning");
|
||||||
|
|
||||||
|
assert_eq!(replication.targets["arn:test:threshold"].after_threshold_count, 1);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn checked_flatten_rejects_dangling_child() {
|
fn checked_flatten_rejects_dangling_child() {
|
||||||
let root_key = hash_path("bucket").key();
|
let root_key = hash_path("bucket").key();
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ pub const SUFFIX_CRC: &str = "crc";
|
|||||||
pub const SUFFIX_TRANSITION_STATUS: &str = "transition-status";
|
pub const SUFFIX_TRANSITION_STATUS: &str = "transition-status";
|
||||||
pub const SUFFIX_TRANSITIONED_OBJECTNAME: &str = "transitioned-object";
|
pub const SUFFIX_TRANSITIONED_OBJECTNAME: &str = "transitioned-object";
|
||||||
pub const SUFFIX_TRANSITIONED_VERSION_ID: &str = "transitioned-versionID";
|
pub const SUFFIX_TRANSITIONED_VERSION_ID: &str = "transitioned-versionID";
|
||||||
|
pub const SUFFIX_TRANSITIONED_VERSION_STATE: &str = "transitioned-version-state";
|
||||||
pub const SUFFIX_TRANSITION_TIER: &str = "transition-tier";
|
pub const SUFFIX_TRANSITION_TIER: &str = "transition-tier";
|
||||||
pub const SUFFIX_TRANSITION_TIER_DESTINATION_ID: &str = "transition-tier-destination-id";
|
pub const SUFFIX_TRANSITION_TIER_DESTINATION_ID: &str = "transition-tier-destination-id";
|
||||||
pub const SUFFIX_RESTORE_OPERATION_ID: &str = "restore-operation-id";
|
pub const SUFFIX_RESTORE_OPERATION_ID: &str = "restore-operation-id";
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ for later deletion.
|
|||||||
- `disk-mutation-body-digest` internode mutating disk RPCs: servers temporarily accept mutating disk RPCs (RenameData, DeleteVersion, DeleteVersions, WriteMetadata, UpdateMetadata, WriteAll, Delete, DeletePaths, RenameFile, RenamePart, DeleteVolume, MakeVolume, MakeVolumes) that carry no signature-bound canonical body digest, so peers from releases that predate body-digest signing remain available during rolling upgrades. Accepted digestless mutations increment the internode body-digest fallback counter; that counter must read zero fleet-wide across a release window before RUSTFS_INTERNODE_RPC_BODY_DIGEST_STRICT is enabled. Because body-bound requests now consume replay-cache nonces on the receiver, deploy the raised RUSTFS_INTERNODE_RPC_REPLAY_CACHE_CAPACITY default fleet-wide before enabling strict mode, and watch the internode replay-cache overflow counter for undersized capacity during the rollout. Remove the digestless fallback after the minimum supported RustFS peer version body-binds every mutating disk RPC.
|
- `disk-mutation-body-digest` internode mutating disk RPCs: servers temporarily accept mutating disk RPCs (RenameData, DeleteVersion, DeleteVersions, WriteMetadata, UpdateMetadata, WriteAll, Delete, DeletePaths, RenameFile, RenamePart, DeleteVolume, MakeVolume, MakeVolumes) that carry no signature-bound canonical body digest, so peers from releases that predate body-digest signing remain available during rolling upgrades. Accepted digestless mutations increment the internode body-digest fallback counter; that counter must read zero fleet-wide across a release window before RUSTFS_INTERNODE_RPC_BODY_DIGEST_STRICT is enabled. Because body-bound requests now consume replay-cache nonces on the receiver, deploy the raised RUSTFS_INTERNODE_RPC_REPLAY_CACHE_CAPACITY default fleet-wide before enabling strict mode, and watch the internode replay-cache overflow counter for undersized capacity during the rollout. Remove the digestless fallback after the minimum supported RustFS peer version body-binds every mutating disk RPC.
|
||||||
- `heal-status-rpc-v1` node heal status capability: new peers treat an unimplemented BackgroundHealStatus RPC as an explicitly incomplete rolling-upgrade response. Remove the fallback after the minimum supported RustFS peer version implements BackgroundHealStatus.
|
- `heal-status-rpc-v1` node heal status capability: new peers treat an unimplemented BackgroundHealStatus RPC as an explicitly incomplete rolling-upgrade response. Remove the fallback after the minimum supported RustFS peer version implements BackgroundHealStatus.
|
||||||
- `backlog-1316` legacy encrypted multipart range seek: the feature remains opt-in until every server that can initiate, write, or complete multipart uploads supports the candidate-to-final marker protocol and uploadId commit lock, and pre-upgrade multipart uploads have drained. Remove the RUSTFS_ENCRYPTED_RANGE_SEEK switch after the minimum supported release does so; keep the quorum marker and malformed-layout full-read guards permanently.
|
- `backlog-1316` legacy encrypted multipart range seek: the feature remains opt-in until every server that can initiate, write, or complete multipart uploads supports the candidate-to-final marker protocol and uploadId commit lock, and pre-upgrade multipart uploads have drained. Remove the RUSTFS_ENCRYPTED_RANGE_SEEK switch after the minimum supported release does so; keep the quorum marker and malformed-layout full-read guards permanently.
|
||||||
|
- `tonic-013-status-render` peer RPC failure classification: internode failures that reach a node only as text (a peer's error_info payload, a status flattened through format!) are classified by matching the rendering of an Unavailable gRPC status. Releases up to 1.0.0-alpha.38 shipped tonic 0.13, which rendered that status as "status: Unavailable, message: ..."; tonic 0.14 renders it as "code: 'The service is currently unavailable', message: ...". Both forms are matched so an older peer's relayed text still marks an unreachable peer offline. Remove the tonic 0.13 form after the minimum supported RustFS peer version ships tonic 0.14 or later.
|
||||||
- `rustfs-5063` pre-beta.9 Local KMS recovery: persisted Local KMS configs from beta.8 and earlier predate the explicit insecure-development flag, and encrypted key files use the legacy SHA-256 KDF. Remove the config fallback after supported upgrades have rewritten or explicitly resaved all pre-beta.9 configs with the development-default field, and remove the legacy KDF after supported upgrades have rewritten all pre-beta.9 Local KMS key files with explicit at-rest protection.
|
- `rustfs-5063` pre-beta.9 Local KMS recovery: persisted Local KMS configs from beta.8 and earlier predate the explicit insecure-development flag, and encrypted key files use the legacy SHA-256 KDF. Remove the config fallback after supported upgrades have rewritten or explicitly resaved all pre-beta.9 configs with the development-default field, and remove the legacy KDF after supported upgrades have rewritten all pre-beta.9 Local KMS key files with explicit at-rest protection.
|
||||||
- `sse-local-dek-json-v1` legacy local SSE DEK decoding: releases before the JSON envelope wrote wrapped DEKs as `base64(nonce):base64(ciphertext)`, so readers retain that decoder while all new writes use the versioned JSON envelope. Remove the colon decoder after the minimum supported direct-upgrade release writes JSON envelopes and migration tooling has rewritten every retained legacy object.
|
- `sse-local-dek-json-v1` legacy local SSE DEK decoding: releases before the JSON envelope wrote wrapped DEKs as `base64(nonce):base64(ciphertext)`, so readers retain that decoder while all new writes use the versioned JSON envelope. Remove the colon decoder after the minimum supported direct-upgrade release writes JSON envelopes and migration tooling has rewritten every retained legacy object.
|
||||||
|
|
||||||
|
|||||||
@@ -763,7 +763,7 @@ async fn enqueue_transitioned_delete_cleanup(
|
|||||||
let _activity_guard = DeleteTailActivityGuard::new(DeleteTailStage::Cleanup);
|
let _activity_guard = DeleteTailActivityGuard::new(DeleteTailStage::Cleanup);
|
||||||
|
|
||||||
let je = if opts.delete_prefix {
|
let je = if opts.delete_prefix {
|
||||||
tier_sweeper::transitioned_force_delete_journal_entry(&existing.transitioned_object)
|
tier_sweeper::transitioned_force_delete_journal_entry(&existing.transitioned_object, existing.transition_version_state)
|
||||||
} else {
|
} else {
|
||||||
let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok());
|
let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok());
|
||||||
tier_sweeper::transitioned_delete_journal_entry(
|
tier_sweeper::transitioned_delete_journal_entry(
|
||||||
@@ -771,6 +771,7 @@ async fn enqueue_transitioned_delete_cleanup(
|
|||||||
opts.versioned,
|
opts.versioned,
|
||||||
opts.version_suspended,
|
opts.version_suspended,
|
||||||
&existing.transitioned_object,
|
&existing.transitioned_object,
|
||||||
|
existing.transition_version_state,
|
||||||
)
|
)
|
||||||
};
|
};
|
||||||
let Some(mut je) = je else {
|
let Some(mut je) = je else {
|
||||||
@@ -9683,7 +9684,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn transitioned_delete_cleanup_persists_identity_bound_and_legacy_journals() {
|
async fn transitioned_delete_cleanup_persists_known_state_and_rejects_unknown_state() {
|
||||||
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
|
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
|
||||||
if current_app_context().is_none() {
|
if current_app_context().is_none() {
|
||||||
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
|
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
|
||||||
@@ -9703,8 +9704,9 @@ mod tests {
|
|||||||
current.transitioned_object.tier = "WARM".to_string();
|
current.transitioned_object.tier = "WARM".to_string();
|
||||||
current.transitioned_object.name = "remote/identity-bound".to_string();
|
current.transitioned_object.name = "remote/identity-bound".to_string();
|
||||||
current.transitioned_object.version_id = "remote-version".to_string();
|
current.transitioned_object.version_id = "remote-version".to_string();
|
||||||
|
current.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
|
||||||
|
|
||||||
let journal_name = |remote_object: &str, backend_identity: Option<[u8; 32]>| {
|
let journal_name = |remote_object: &str, backend_identity: Option<[u8; 32]>, version_id_exact: bool| {
|
||||||
use sha2::{Digest, Sha256};
|
use sha2::{Digest, Sha256};
|
||||||
|
|
||||||
let mut hasher = Sha256::new();
|
let mut hasher = Sha256::new();
|
||||||
@@ -9717,6 +9719,10 @@ mod tests {
|
|||||||
hasher.update([0]);
|
hasher.update([0]);
|
||||||
hasher.update(backend_identity);
|
hasher.update(backend_identity);
|
||||||
}
|
}
|
||||||
|
if version_id_exact {
|
||||||
|
hasher.update([0]);
|
||||||
|
hasher.update(b"exact-version-id");
|
||||||
|
}
|
||||||
format!("ilm/tier-delete-journal/{}.json", rustfs_utils::crypto::hex(hasher.finalize().as_slice()))
|
format!("ilm/tier-delete-journal/{}.json", rustfs_utils::crypto::hex(hasher.finalize().as_slice()))
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -9726,7 +9732,7 @@ mod tests {
|
|||||||
let mut identity_bound = store
|
let mut identity_bound = store
|
||||||
.get_object_reader(
|
.get_object_reader(
|
||||||
".rustfs.sys",
|
".rustfs.sys",
|
||||||
&journal_name("remote/identity-bound", Some(identity)),
|
&journal_name("remote/identity-bound", Some(identity), true),
|
||||||
None,
|
None,
|
||||||
http::HeaderMap::new(),
|
http::HeaderMap::new(),
|
||||||
&ObjectOptions::default(),
|
&ObjectOptions::default(),
|
||||||
@@ -9739,11 +9745,14 @@ mod tests {
|
|||||||
.expect("identity-bound journal body should be readable");
|
.expect("identity-bound journal body should be readable");
|
||||||
let identity_bound: serde_json::Value =
|
let identity_bound: serde_json::Value =
|
||||||
serde_json::from_slice(&identity_bound_data).expect("identity-bound journal should decode as JSON");
|
serde_json::from_slice(&identity_bound_data).expect("identity-bound journal should decode as JSON");
|
||||||
assert_eq!(identity_bound["version"], serde_json::json!(2));
|
assert_eq!(identity_bound["version"], serde_json::json!(4));
|
||||||
assert_eq!(identity_bound["backend_identity"], serde_json::json!(identity));
|
assert_eq!(identity_bound["backend_identity"], serde_json::json!(identity));
|
||||||
|
assert_eq!(identity_bound["version_id_exact"], serde_json::json!(true));
|
||||||
|
assert_eq!(identity_bound["version_state"], serde_json::json!("exact"));
|
||||||
|
|
||||||
current.user_defined = Arc::new(HashMap::new());
|
current.user_defined = Arc::new(HashMap::new());
|
||||||
current.transitioned_object.name = "remote/legacy".to_string();
|
current.transitioned_object.name = "remote/legacy".to_string();
|
||||||
|
current.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown;
|
||||||
enqueue_transitioned_delete_cleanup(
|
enqueue_transitioned_delete_cleanup(
|
||||||
store.clone(),
|
store.clone(),
|
||||||
"bucket",
|
"bucket",
|
||||||
@@ -9755,24 +9764,31 @@ mod tests {
|
|||||||
Some(¤t),
|
Some(¤t),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("legacy force-delete cleanup should persist a fail-closed v1 journal");
|
.expect("unknown force-delete cleanup should fail closed without a journal");
|
||||||
let mut legacy = store
|
let legacy_err = match store
|
||||||
.get_object_reader(
|
.get_object_reader(
|
||||||
".rustfs.sys",
|
".rustfs.sys",
|
||||||
&journal_name("remote/legacy", None),
|
&journal_name("remote/legacy", None, false),
|
||||||
None,
|
None,
|
||||||
http::HeaderMap::new(),
|
http::HeaderMap::new(),
|
||||||
&ObjectOptions::default(),
|
&ObjectOptions::default(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.expect("legacy journal should be readable");
|
{
|
||||||
let mut legacy_data = Vec::new();
|
Ok(_) => panic!("unknown remote version state must not persist a delete journal"),
|
||||||
tokio::io::AsyncReadExt::read_to_end(&mut legacy.stream, &mut legacy_data)
|
Err(err) => err,
|
||||||
.await
|
};
|
||||||
.expect("legacy journal body should be readable");
|
assert!(
|
||||||
let legacy: serde_json::Value = serde_json::from_slice(&legacy_data).expect("legacy journal should decode as JSON");
|
matches!(
|
||||||
assert_eq!(legacy["version"], serde_json::json!(1));
|
&legacy_err,
|
||||||
assert_eq!(legacy["backend_identity"], serde_json::Value::Null);
|
StorageError::FileNotFound
|
||||||
|
| StorageError::ObjectNotFound(_, _)
|
||||||
|
| StorageError::FileVersionNotFound
|
||||||
|
| StorageError::VersionNotFound(_, _, _)
|
||||||
|
| StorageError::VolumeNotFound
|
||||||
|
),
|
||||||
|
"unknown remote version state must leave no journal, got {legacy_err:?}"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn put_real_cold_fill_object(store: &Arc<ECStore>, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo {
|
async fn put_real_cold_fill_object(store: &Arc<ECStore>, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo {
|
||||||
|
|||||||
@@ -409,20 +409,24 @@ pub(crate) mod bucket {
|
|||||||
versioned: bool,
|
versioned: bool,
|
||||||
suspended: bool,
|
suspended: bool,
|
||||||
transitioned: &super::super::super::storage_contracts::TransitionedObject,
|
transitioned: &super::super::super::storage_contracts::TransitionedObject,
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
) -> Option<Jentry> {
|
) -> Option<Jentry> {
|
||||||
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry(
|
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry(
|
||||||
version_id,
|
version_id,
|
||||||
versioned,
|
versioned,
|
||||||
suspended,
|
suspended,
|
||||||
transitioned,
|
transitioned,
|
||||||
|
transition_version_state,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn transitioned_force_delete_journal_entry(
|
pub(crate) fn transitioned_force_delete_journal_entry(
|
||||||
transitioned: &super::super::super::storage_contracts::TransitionedObject,
|
transitioned: &super::super::super::storage_contracts::TransitionedObject,
|
||||||
|
transition_version_state: rustfs_filemeta::TransitionVersionState,
|
||||||
) -> Option<Jentry> {
|
) -> Option<Jentry> {
|
||||||
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry(
|
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry(
|
||||||
transitioned,
|
transitioned,
|
||||||
|
transition_version_state,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,14 +13,7 @@ fi
|
|||||||
classify_source_layer() {
|
classify_source_layer() {
|
||||||
local file="$1"
|
local file="$1"
|
||||||
|
|
||||||
if [[ "$file" == rustfs/src/server/* ]] ||
|
if [[ "$file" == rustfs/src/app/* ]]; then
|
||||||
[[ "$file" == rustfs/src/startup_*.rs ]] ||
|
|
||||||
[[ "$file" == rustfs/src/init.rs ]] ||
|
|
||||||
[[ "$file" == rustfs/src/main.rs ]] ||
|
|
||||||
[[ "$file" == rustfs/src/lib.rs ]] ||
|
|
||||||
[[ "$file" == rustfs/src/embedded.rs ]]; then
|
|
||||||
printf 'composition'
|
|
||||||
elif [[ "$file" == rustfs/src/app/* ]]; then
|
|
||||||
printf 'app'
|
printf 'app'
|
||||||
elif [[ "$file" == rustfs/src/admin/* ]] || [[ "$file" == rustfs/src/storage/ecfs.rs ]] || [[ "$file" == rustfs/src/storage/s3_api/* ]]; then
|
elif [[ "$file" == rustfs/src/admin/* ]] || [[ "$file" == rustfs/src/storage/ecfs.rs ]] || [[ "$file" == rustfs/src/storage/s3_api/* ]]; then
|
||||||
printf 'interface'
|
printf 'interface'
|
||||||
@@ -37,14 +30,6 @@ classify_target_layer() {
|
|||||||
local storage_path
|
local storage_path
|
||||||
|
|
||||||
case "$root" in
|
case "$root" in
|
||||||
init | main | lib | embedded | startup_*)
|
|
||||||
printf 'composition'
|
|
||||||
;;
|
|
||||||
server)
|
|
||||||
# Server files are composition roots when they import lower layers, but
|
|
||||||
# their exported HTTP contracts belong to the interface boundary.
|
|
||||||
printf 'interface'
|
|
||||||
;;
|
|
||||||
app)
|
app)
|
||||||
printf 'app'
|
printf 'app'
|
||||||
;;
|
;;
|
||||||
@@ -68,9 +53,6 @@ classify_target_layer() {
|
|||||||
|
|
||||||
layer_rank() {
|
layer_rank() {
|
||||||
case "$1" in
|
case "$1" in
|
||||||
composition)
|
|
||||||
printf '4'
|
|
||||||
;;
|
|
||||||
interface)
|
interface)
|
||||||
printf '3'
|
printf '3'
|
||||||
;;
|
;;
|
||||||
@@ -86,48 +68,6 @@ layer_rank() {
|
|||||||
esac
|
esac
|
||||||
}
|
}
|
||||||
|
|
||||||
is_reverse_dependency() {
|
|
||||||
local source_rank target_rank
|
|
||||||
|
|
||||||
source_rank="$(layer_rank "$1")"
|
|
||||||
target_rank="$(layer_rank "$2")"
|
|
||||||
(( source_rank < target_rank ))
|
|
||||||
}
|
|
||||||
|
|
||||||
assert_dependency_direction() {
|
|
||||||
local expected="$1"
|
|
||||||
local source="$2"
|
|
||||||
local target="$3"
|
|
||||||
local actual='allowed'
|
|
||||||
|
|
||||||
if is_reverse_dependency "$source" "$target"; then
|
|
||||||
actual='reverse'
|
|
||||||
fi
|
|
||||||
if [[ "$actual" != "$expected" ]]; then
|
|
||||||
printf 'Layer dependency guard self-test failed: %s -> %s (expected %s, got %s)\n' \
|
|
||||||
"$source" "$target" "$expected" "$actual" >&2
|
|
||||||
exit 1
|
|
||||||
fi
|
|
||||||
}
|
|
||||||
|
|
||||||
run_layer_model_self_tests() {
|
|
||||||
local server_source app_source storage_source server_target admin_target app_target
|
|
||||||
|
|
||||||
server_source="$(classify_source_layer rustfs/src/server/http.rs)"
|
|
||||||
app_source="$(classify_source_layer rustfs/src/app/bucket_usecase.rs)"
|
|
||||||
storage_source="$(classify_source_layer rustfs/src/storage/rpc/node_service.rs)"
|
|
||||||
server_target="$(classify_target_layer server::http)"
|
|
||||||
admin_target="$(classify_target_layer admin::router)"
|
|
||||||
app_target="$(classify_target_layer app::bucket_usecase)"
|
|
||||||
|
|
||||||
assert_dependency_direction 'allowed' "$server_source" "$admin_target"
|
|
||||||
assert_dependency_direction 'allowed' "$server_source" "$app_target"
|
|
||||||
assert_dependency_direction 'reverse' "$app_source" "$server_target"
|
|
||||||
assert_dependency_direction 'reverse' "$storage_source" "$server_target"
|
|
||||||
assert_dependency_direction 'reverse' "$app_source" "$admin_target"
|
|
||||||
assert_dependency_direction 'reverse' "$storage_source" "$admin_target"
|
|
||||||
}
|
|
||||||
|
|
||||||
normalize_import_group_item() {
|
normalize_import_group_item() {
|
||||||
local prefix="$1"
|
local prefix="$1"
|
||||||
local item="$2"
|
local item="$2"
|
||||||
@@ -242,11 +182,6 @@ normalize_import_path() {
|
|||||||
|
|
||||||
emit_crate_use_statements() {
|
emit_crate_use_statements() {
|
||||||
(cd "$ROOT_DIR" && rg --files -g '*.rs' rustfs/src | while IFS= read -r file; do
|
(cd "$ROOT_DIR" && rg --files -g '*.rs' rustfs/src | while IFS= read -r file; do
|
||||||
# The guard is file-scoped: dedicated test modules are excluded, while
|
|
||||||
# inline #[cfg(test)] imports remain subject to the source file's layer.
|
|
||||||
if [[ "$file" == *_test.rs ]] || [[ "$file" == */tests/* ]]; then
|
|
||||||
continue
|
|
||||||
fi
|
|
||||||
perl -0777 -ne '
|
perl -0777 -ne '
|
||||||
while (/\buse\s+crate::.*?;/sg) {
|
while (/\buse\s+crate::.*?;/sg) {
|
||||||
my $statement = $&;
|
my $statement = $&;
|
||||||
@@ -259,26 +194,6 @@ emit_crate_use_statements() {
|
|||||||
done)
|
done)
|
||||||
}
|
}
|
||||||
|
|
||||||
write_baseline_file() {
|
|
||||||
local entries="$1"
|
|
||||||
|
|
||||||
cat >"$BASELINE_FILE" <<'EOF'
|
|
||||||
# Layer dependency baseline for the rustfs binary crate.
|
|
||||||
#
|
|
||||||
# The guard models production imports as:
|
|
||||||
# composition -> interface -> app -> infra
|
|
||||||
#
|
|
||||||
# Canonical dependency entry:
|
|
||||||
# dep|source_file|source_layer->target_layer|crate::imported_symbol
|
|
||||||
#
|
|
||||||
# Canonical conceptual cycle entry:
|
|
||||||
# cycle|left_layer<->right_layer
|
|
||||||
EOF
|
|
||||||
cat "$entries" >>"$BASELINE_FILE"
|
|
||||||
}
|
|
||||||
|
|
||||||
run_layer_model_self_tests
|
|
||||||
|
|
||||||
normalize_baseline_file() {
|
normalize_baseline_file() {
|
||||||
local input="$1"
|
local input="$1"
|
||||||
local output="$2"
|
local output="$2"
|
||||||
@@ -350,7 +265,10 @@ while IFS= read -r line; do
|
|||||||
printf '%s->%s\n' "$source_layer" "$target_layer" >>"$EDGES_RAW"
|
printf '%s->%s\n' "$source_layer" "$target_layer" >>"$EDGES_RAW"
|
||||||
fi
|
fi
|
||||||
|
|
||||||
if is_reverse_dependency "$source_layer" "$target_layer"; then
|
source_rank="$(layer_rank "$source_layer")"
|
||||||
|
target_rank="$(layer_rank "$target_layer")"
|
||||||
|
|
||||||
|
if (( source_rank < target_rank )); then
|
||||||
printf 'dep|%s|%s->%s|crate::%s\n' "$file" "$source_layer" "$target_layer" "$import_path" >>"$VIOLATIONS_RAW"
|
printf 'dep|%s|%s->%s|crate::%s\n' "$file" "$source_layer" "$target_layer" "$import_path" >>"$VIOLATIONS_RAW"
|
||||||
fi
|
fi
|
||||||
done < <(normalize_import_path "$text")
|
done < <(normalize_import_path "$text")
|
||||||
@@ -374,7 +292,7 @@ done <"${TMP_DIR}/edges_sorted.txt" | sort -u >"${TMP_DIR}/cycles_sorted.txt"
|
|||||||
cat "${TMP_DIR}/violations_sorted.txt" "${TMP_DIR}/cycles_sorted.txt" | sort -u >"$CURRENT_BASELINE"
|
cat "${TMP_DIR}/violations_sorted.txt" "${TMP_DIR}/cycles_sorted.txt" | sort -u >"$CURRENT_BASELINE"
|
||||||
|
|
||||||
if [[ "$MODE" == "update" ]]; then
|
if [[ "$MODE" == "update" ]]; then
|
||||||
write_baseline_file "$CURRENT_BASELINE"
|
cp "$CURRENT_BASELINE" "$BASELINE_FILE"
|
||||||
echo "Updated baseline: $BASELINE_FILE"
|
echo "Updated baseline: $BASELINE_FILE"
|
||||||
exit 0
|
exit 0
|
||||||
fi
|
fi
|
||||||
|
|||||||
@@ -1,44 +1,35 @@
|
|||||||
# Layer dependency baseline for the rustfs binary crate.
|
# Layer dependency baseline for the rustfs binary crate.
|
||||||
#
|
#
|
||||||
# The guard models production imports as:
|
# These are intra-crate module references within rustfs/ that cross the
|
||||||
# composition -> interface -> app -> infra
|
# conceptual layer boundaries (app, infra/server, interface/admin).
|
||||||
|
# Since they live inside one Cargo crate, Rust doesn't enforce separation.
|
||||||
|
# The list is maintained for architectural awareness during code review.
|
||||||
#
|
#
|
||||||
# Canonical dependency entry:
|
# Format for dependency entries:
|
||||||
# dep|source_file|source_layer->target_layer|crate::imported_symbol
|
# status|source_file|direction|imported_symbol|reason
|
||||||
#
|
#
|
||||||
# Canonical conceptual cycle entry:
|
# Format for conceptual cycles:
|
||||||
# cycle|left_layer<->right_layer
|
# status|cycle|direction_pair|reason
|
||||||
|
#
|
||||||
|
# Status:
|
||||||
|
# accepted - reviewed and intentionally allowed
|
||||||
|
# todo - should be resolved in a future refactor
|
||||||
|
|
||||||
cycle|app<->infra
|
cycle|app<->infra
|
||||||
cycle|app<->interface
|
cycle|app<->interface
|
||||||
cycle|composition<->infra
|
|
||||||
cycle|composition<->interface
|
|
||||||
cycle|infra<->interface
|
cycle|infra<->interface
|
||||||
dep|rustfs/src/admin/handlers/scanner.rs|interface->composition|crate::startup_background::ENV_SCANNER_ENABLED
|
|
||||||
dep|rustfs/src/admin/handlers/scanner.rs|interface->composition|crate::startup_background::scanner_enabled_from_env
|
|
||||||
dep|rustfs/src/app/admin_usecase.rs|app->interface|crate::server::DependencyReadiness
|
|
||||||
dep|rustfs/src/app/admin_usecase.rs|app->interface|crate::server::collect_dependency_readiness_report
|
|
||||||
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_bucket_meta_hook
|
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_bucket_meta_hook
|
||||||
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_delete_bucket_hook
|
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_delete_bucket_hook
|
||||||
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_make_bucket_hook
|
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::admin::handlers::site_replication::site_replication_make_bucket_hook
|
||||||
dep|rustfs/src/app/bucket_usecase.rs|app->interface|crate::server::RemoteAddr
|
dep|rustfs/src/init.rs|infra->interface|crate::admin
|
||||||
dep|rustfs/src/app/object_usecase.rs|app->interface|crate::server::convert_ecstore_object_info
|
|
||||||
dep|rustfs/src/cluster_snapshot.rs|infra->interface|crate::server::DependencyReadiness
|
|
||||||
dep|rustfs/src/cluster_snapshot.rs|infra->interface|crate::server::DependencyReadinessReport
|
|
||||||
dep|rustfs/src/cluster_snapshot.rs|infra->interface|crate::server::ReadinessDegradedReason
|
|
||||||
dep|rustfs/src/cluster_snapshot.rs|infra->interface|crate::server::snapshot_dependency_readiness_report
|
|
||||||
dep|rustfs/src/runtime_sources.rs|infra->app|crate::app::context
|
dep|rustfs/src/runtime_sources.rs|infra->app|crate::app::context
|
||||||
dep|rustfs/src/storage/access.rs|infra->interface|crate::server::RemoteAddr
|
dep|rustfs/src/server/http.rs|infra->interface|crate::admin
|
||||||
dep|rustfs/src/storage/ecfs_extend.rs|infra->interface|crate::server::cors
|
dep|rustfs/src/server/layer.rs|infra->interface|crate::admin::console::is_console_path
|
||||||
dep|rustfs/src/storage/ecfs_extend.rs|infra->interface|crate::storage::ecfs::ListObjectUnorderedQuery
|
dep|rustfs/src/storage/ecfs_extend.rs|infra->interface|crate::storage::ecfs::ListObjectUnorderedQuery
|
||||||
dep|rustfs/src/storage/helper.rs|infra->interface|crate::server::convert_ecstore_object_info
|
dep|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::FS
|
||||||
dep|rustfs/src/storage/helper.rs|infra->interface|crate::server::is_audit_module_enabled
|
dep|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::validate_object_lock_configuration_input
|
||||||
dep|rustfs/src/storage/helper.rs|infra->interface|crate::server::is_notify_module_enabled
|
dep|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::s3_api::common::rustfs_initiator
|
||||||
dep|rustfs/src/storage/helper.rs|infra->interface|crate::server::refresh_audit_module_enabled
|
dep|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::s3_api::common::rustfs_owner
|
||||||
dep|rustfs/src/storage/helper.rs|infra->interface|crate::server::refresh_notify_module_enabled
|
|
||||||
dep|rustfs/src/storage/rpc/http_service.rs|infra->interface|crate::server::RPC_PREFIX
|
|
||||||
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::config::reload_dynamic_config_runtime_state
|
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::config::reload_dynamic_config_runtime_state
|
||||||
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::config::reload_runtime_config_snapshot
|
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::config::reload_runtime_config_snapshot
|
||||||
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::site_replication::reload_site_replication_runtime_state
|
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::admin::service::site_replication::reload_site_replication_runtime_state
|
||||||
dep|rustfs/src/storage/rpc/node_service.rs|infra->interface|crate::server::MODULE_SWITCHES_SIGNAL_SUBSYSTEM
|
|
||||||
dep|rustfs/src/storage/rpc/node_service/heal.rs|infra->composition|crate::startup_background::heal_enabled_from_env
|
|
||||||
dep|rustfs/src/storage/rpc/node_service/heal.rs|infra->composition|crate::startup_background::scanner_enabled_from_env
|
|
||||||
|
|||||||
Reference in New Issue
Block a user