Compare commits

..

32 Commits

Author SHA1 Message Date
houseme d603ea3922 chore: merge main into scanner scope entry oracles
Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-06 11:08:55 +08:00
houseme 4068f8440c test(scanner): verify default scoped entry fallback walks
Exercise storage-derived scope resolution through real per-source walkers,
including invalid baselines, dirty overflow and changed bucket inventory.
Distinguish same-cycle Current retries from unbound cold-bucket reuse, and
assert candidate delivery preserves authoritative root and dirty state.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-06 04:05:37 +08:00
Zhengchao An 982ed2888f Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 03:32:00 +08:00
houseme eec6d3a76d Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 02:29:11 +08:00
houseme bbca907781 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 01:50:19 +08:00
houseme 51982be687 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-06 01:33:21 +08:00
houseme b8f0119320 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 23:53:05 +08:00
houseme 20ca988e85 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 23:33:13 +08:00
houseme 045770d09d fix(scanner): satisfy cache prefix sort lint
Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 22:34:20 +08:00
houseme 278cbaefd5 Merge branch 'main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 22:16:44 +08:00
houseme 982bebf6a5 Merge remote-tracking branch 'origin/main' into houseme/fix/scanner-heal-v2-scanner-integration
Resolved scanner cache publication conflicts after the main branch added
execution identity fencing and post-lease activity proof coverage.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 22:10:39 +08:00
houseme 1aafe477ec test(scanner): verify joint checkpoint coverage metadata
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 21:05:13 +08:00
houseme c33a42c5ab test(scanner): supply explicit coverage in publication fixtures
Keep the confirmed-empty namespace fixture authoritative under the required
coverage contract and qualify the bucket cache metadata test type.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:55:19 +08:00
houseme e3fd1ae19a fix(scanner): fence same-cycle caches with full activity coverage
Keep structural baseline identity separate from the full activity coverage
required by bucket admission and set publication. Require complete set
coverage proofs while retaining revision CAS and epoch regression checks.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:55:19 +08:00
houseme 6c1d401215 test(scanner): reproduce same-cycle dirty aggregate replay
Cover a Normal-to-Normal retry with a new dirty bucket generation after
bucket persistence and root delivery failure.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:53:20 +08:00
houseme c80ff445c0 fix(scanner): fence set snapshot reuse with the scan work proof
Prevent same-cycle set publication from replacing freshly scanned maintenance
results with an older Normal aggregate. Recognize uniform completed
maintenance baselines when planning later ordinary dirty-bucket work.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:53:20 +08:00
houseme 3d4c05ddf9 fix(scanner): bind bucket cache reuse to scan work requirements
Carry stable scan mode and full-maintenance requirements in the existing
opaque bucket digest before local and remote cache admission. Different
requirements cannot replay a same-cycle Normal cache after root delivery
failure; matching requirements remain reusable for the same intent.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:52:15 +08:00
houseme f9dd0388fa fix(scanner): refresh scope safety independently of idle backoff
Inspect maintenance on multi-disk startup and refresh changed or failed
evidence even when explicit bitrot configuration disables idle backoff.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:50:13 +08:00
houseme 527860d71e fix(scanner): keep maintenance cycles outside dirty bucket scopes
Force complete bucket scope for deep scans and scheduled maintenance while
preserving the existing planner for verified ordinary dirty work.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:50:12 +08:00
houseme 7d2c120073 test(scanner): use valid modification times in checkpoint fixtures
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme 489c3276a4 fix(scanner): verify coverage receipts and scan strength
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme b09a4f3706 fix(scanner): keep stable snapshot rescan behavior
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme a0c32e7c6b fix(scanner): retain scoped partial coverage across dirty plans
Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:49:50 +08:00
houseme d5a1530409 fix(scanner): require complete publication coverage
Refs rustfs/backlog#2261 and rustfs/backlog#2240.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 20:48:01 +08:00
houseme e7b0349a4f Merge remote-tracking branch 'origin/main' into houseme/fix/scanner-heal-v2-scanner-integration 2026-09-05 20:43:53 +08:00
houseme bdbdca07c8 fix(deps): preserve supported hotpath focus expressions
Keep the profiler runtime before its regex-lite compatibility regression.
Track the opt-in validation required to remove this constraint in backlog.

Refs rustfs/backlog#2302.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:50:46 +08:00
houseme 53efaa2b8f chore(deps): refresh profiling dependencies for the next batch
Update hotpath and its macro crate to the compatible patch release before
the next dependency-ready implementation tasks.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 19:10:54 +08:00
houseme ec9672a397 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b4-base 2026-09-05 18:54:50 +08:00
houseme cee35f7e54 Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:43:37 +08:00
houseme ef7e7afd8c Merge remote-tracking branch 'origin/main' into houseme/chore/scanner-heal-v2-b3-base 2026-09-05 16:35:05 +08:00
houseme 652ebb12c6 fix(ecstore): remove duplicate local rename implementation
Keep the canonical commit module after concurrent storage changes merged.
The control-write and rollback changes are already present there.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:28:01 +08:00
houseme 38c03d9d5d chore(deps): refresh scanner heal batch dependency baseline
Regenerate compatible lockfile selections before the next implementation
batch. Cargo upgrade leaves direct requirements unchanged.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-05 16:22:13 +08:00
5 changed files with 305 additions and 272 deletions
+11 -21
View File
@@ -22,13 +22,11 @@
# Upgrade cases download the same pinned previous release as e2e-upgrade.yml.
#
# Isolated pool filesystems: expand/decommission/rebalance cases require
# independent `statfs` capacity. This job runs on GitHub-hosted
# `ubuntu-latest` because the self-hosted `sm-standard-4` ARC pods cannot
# create filesystems: `mount -o loop` fails with ENOENT (no
# `/dev/loop-control`), and `mount -t tmpfs` fails with "cannot mount tmpfs
# read-only" (no `CAP_SYS_ADMIN` in the initial namespace). The same reason
# `uring-integration` and `e2e-s3tests.yml` left that label. The prepare
# step mounts four 1 GiB tmpfs instances and exports `RUSTFS_E2E_POOL_ROOTS`.
# independent `statfs` capacity. `sm-standard-4` is an ARC pod
# (`scripts/ci/check_runner_ephemerality.sh`) and usually has no
# `/dev/loop-control`, so `mount -o loop` fails with ENOENT ("mount failed:
# No such file or directory"). The prepare step therefore mounts four 1 GiB
# tmpfs instances and exports them as `RUSTFS_E2E_POOL_ROOTS`.
name: e2e-distributed
@@ -78,9 +76,7 @@ concurrency:
jobs:
distributed:
name: Distributed 4-node 4-disk e2e
# GitHub-hosted VM: loop and tmpfs mounts work here. sm-standard-4 is an
# ARC pod and rejects both (`mount -o loop` ENOENT, tmpfs "read-only").
runs-on: ubuntu-latest
runs-on: sm-standard-4
timeout-minutes: 180
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
@@ -101,9 +97,7 @@ jobs:
uses: ./.github/actions/setup
with:
rust-version: stable
# Dedicated key: ubuntu-latest and sm-standard-4 share runner.os, so
# a shared key would mix VM and ARC pod target/ artifacts.
cache-shared-key: ci-e2e-distributed-hosted
cache-shared-key: ci-e2e-distributed
cache-save-if: ${{ github.ref == 'refs/heads/main' }}
install-build-packaging-tools: 'false'
@@ -116,14 +110,10 @@ jobs:
for pool in 0 1 2 3; do
mountpoint="${mount_base}/pool-${pool}"
mkdir -p "${mountpoint}"
# Sized tmpfs reports a distinct st_dev and independent 1G
# statfs capacity. Requires a VM runner (ubuntu-latest).
if ! sudo mount -t tmpfs -o size=1G,nosuid,nodev,mode=1777 tmpfs "${mountpoint}"; then
echo "tmpfs mount failed on $(uname -a)" >&2
findmnt || true
grep Cap /proc/self/status || true
exit 1
fi
# sm-standard-4 is an ARC pod without usable loop devices, so
# `mount -o loop` fails with ENOENT. Sized tmpfs still reports a
# distinct st_dev and independent 1G statfs capacity.
sudo mount -t tmpfs -o size=1G,nosuid,nodev,mode=1777 tmpfs "${mountpoint}"
sudo chmod 1777 "${mountpoint}"
roots+=("${mountpoint}")
done
+2
View File
@@ -39,6 +39,8 @@ use temp_env::with_var;
use time::OffsetDateTime;
use uuid::Uuid;
mod scoped_entry_fallback;
#[derive(Clone)]
struct FixedWorkloadProvider {
snapshot: WorkloadAdmissionRegistrySnapshot,
@@ -0,0 +1,289 @@
// Copyright 2026 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::*;
use crate::data_usage_define::{DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision};
use crate::storage_api::owner::EcstoreDiskAPI;
type DriveIdentities = HashMap<String, (Uuid, DataUsageCacheSource)>;
type WalkCounts = HashMap<(String, String, String), u64>;
async fn drive_identities(store: &ECStore) -> DriveIdentities {
let mut identities = HashMap::new();
let mut ids = HashSet::new();
for set in store.all_set_disks() {
let source = DataUsageCacheSource::new(set.pool_index, set.set_index);
for disk in scanner_set_disk_inventory(set.as_ref()).await {
let id = EcstoreDiskAPI::get_disk_id(disk.as_ref())
.await
.expect("fixture disk identity should be readable")
.expect("fixture disk must have a durable identity");
assert!(!id.is_nil());
assert!(ids.insert(id), "fixture disk identities must be unique");
let path = crate::ScannerDiskExt::path(disk.as_ref()).to_string_lossy().into_owned();
assert!(identities.insert(path, (id, source)).is_none());
}
}
assert_eq!(identities.len(), 8);
identities
}
fn walk_counts(drives: &DriveIdentities) -> WalkCounts {
rustfs_scanner_metrics::metrics::global_metrics()
.scanner_runtime_details_report()
.bucket_drive_results
.into_iter()
.filter(|result| drives.contains_key(&result.drive))
.map(|result| ((result.bucket, result.drive, result.result), result.count))
.collect()
}
async fn put_and_settle(store: &ECStore, bucket: &str, object: &str) {
let set = &store.pools[0].disk_set[0];
let mut reader = ScannerPutObjReader::from_vec(b"object".to_vec());
set.put_object(bucket, object, &mut reader, &ScannerObjectOptions::default())
.await
.expect("fixture object should persist");
let lock = set.new_ns_lock(bucket, object).await.expect("fixture namespace lock");
let _settled = lock
.get_write_lock(Duration::from_secs(30))
.await
.expect("quorum-ACK rename tail must settle before taking the activity baseline");
}
async fn create_bucket(store: &ECStore, bucket: &str) {
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("fixture bucket should be created");
put_and_settle(store, bucket, "initial").await;
}
async fn persist_baseline(store: &Arc<ECStore>, baseline: &DataUsageInfo) {
let mut baseline = baseline.clone();
baseline.usage_snapshot_converged = Some(true);
crate::save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&baseline).expect("baseline should encode"),
)
.await
.expect("fixture baseline should persist");
}
// Every invocation uses the production default scope. The expected walker set
// comes from storage's per-source inventory, not the resolver's selected names.
async fn run_entry(store: &Arc<ECStore>, cycle: u64, selected: Option<&str>, expect_walks: bool) -> DataUsageInfo {
let drives = drive_identities(store).await;
let inventory = store
.list_bucket_for_scanner(&BucketOptions::default())
.await
.expect("fixture inventory should be complete");
assert!(inventory.topology_complete);
let expected_walks = if expect_walks {
inventory
.set_buckets
.into_iter()
.flat_map(|set| {
let source = DataUsageCacheSource::new(set.pool_index, set.set_index);
set.buckets.into_iter().map(move |bucket| ((source, bucket.name), 1_u64))
})
.collect::<HashMap<_, _>>()
} else {
HashMap::new()
};
let root_before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("root baseline should be readable");
let dirty_before = dirty_usage_buckets_for_tests();
let generation_before = dirty_usage_generation();
let before = walk_counts(&drives);
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
let (updates, mut receiver) = mpsc::channel(1);
let (observer, observed) = tokio::sync::oneshot::channel();
let result = tokio::time::timeout(
Duration::from_secs(30),
nsscanner_with_storage_status_scoped(
store.as_ref(),
ScannerCycleRequest {
ctx,
budget,
updates,
want_cycle: cycle,
leader_epoch: 11,
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: root_before.0.clone().map(Bytes::from),
requires_full_scan: false,
resolved_scope_observer: Some(observer),
},
),
)
.await
.expect("entry cycle should finish within the fixture deadline")
.expect("entry cycle should succeed");
assert_eq!(result.status, ScannerCycleStatus::Complete);
let scope = observed.await.expect("production resolver should report its decision");
assert_eq!(
scope.selected_buckets.as_deref(),
selected.map(|name| HashSet::from([name.to_string()])).as_ref()
);
let usage = receiver.recv().await.expect("one candidate should be delivered");
assert!(receiver.recv().await.is_none(), "there must be exactly one terminal candidate");
assert!(usage.usage_snapshot_complete);
assert!(!usage.usage_snapshot_partial);
assert_eq!(usage.scanner_cycle, Some(cycle));
assert_eq!(
drive_identities(store).await,
drives,
"drive identities must not change during the oracle"
);
let after = walk_counts(&drives);
let mut actual = HashMap::new();
for key in before.keys() {
assert!(after.contains_key(key), "metrics eviction would invalidate this exact-delta oracle");
}
for ((bucket, drive, outcome), count) in after {
let previous = before
.get(&(bucket.clone(), drive.clone(), outcome.clone()))
.copied()
.unwrap_or(0);
let delta = count.checked_sub(previous).expect("fixture counters must not reset");
if delta > 0 {
assert_eq!(outcome, "success", "no error or partial walker is expected");
*actual.entry((drives[&drive].1, bucket)).or_insert(0_u64) += delta;
}
}
assert_eq!(
actual, expected_walks,
"each listed source/bucket must have exactly the expected real walks"
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("root after scan"),
root_before,
"producing a candidate must not replace the coordinator-owned root baseline"
);
assert_eq!(dirty_usage_generation(), generation_before);
assert!(
dirty_usage_buckets_for_tests() == dirty_before,
"candidate delivery must not ACK pending dirty buckets"
);
usage
}
#[tokio::test]
#[serial]
async fn scoped_entry_fallback_distinguishes_planned_scope_from_real_cold_walks() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
let cold = format!("cold-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
create_bucket(&store, &cold).await;
record_dirty_usage_bucket(&hot);
let baseline = run_entry(&store, 1, None, true).await;
persist_baseline(&store, &baseline).await;
// A same-intent, same-cycle Current cache is a retry, not proof that a
// later cycle may reuse unselected buckets without durable incarnation.
run_entry(&store, 1, Some(&hot), false).await;
let usage = run_entry(&store, 2, Some(&hot), true).await;
assert_eq!(usage.buckets_usage[&hot].objects_count, 1);
assert_eq!(usage.buckets_usage[&cold].objects_count, 1);
assert_eq!(usage.objects_total_count, 2);
clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_entry_fallback_rejects_invalid_persisted_baseline_at_the_walker() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
let cold = format!("cold-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
create_bucket(&store, &cold).await;
record_dirty_usage_bucket(&hot);
// The first real scan is also the missing persisted-baseline case.
let baseline = run_entry(&store, 1, None, true).await;
for (index, kind) in [
"malformed",
"unconverged",
"missing-set",
"wrong-source",
"mixed-plan",
"wrong-epoch",
]
.into_iter()
.enumerate()
{
let mut candidate = baseline.clone();
candidate.usage_snapshot_converged = Some(true);
match kind {
"unconverged" => candidate.usage_snapshot_converged = Some(false),
"missing-set" => {
candidate.usage_snapshot_set_states.pop();
}
"wrong-source" => candidate.usage_snapshot_set_states[0].set_index = 99,
"mixed-plan" => candidate.usage_snapshot_set_states[1].scan_plan_digest = Some([0xA5; 32]),
"wrong-epoch" => candidate.usage_snapshot_set_states[0].scanner_epoch = Some(10),
"malformed" => {}
_ => unreachable!(),
}
let bytes = if kind == "malformed" {
b"{broken".to_vec()
} else {
serde_json::to_vec(&candidate).expect("candidate JSON")
};
crate::save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes)
.await
.expect("negative baseline should persist");
let usage = run_entry(&store, u64::try_from(index).expect("fixture cycle index should fit") + 2, None, true).await;
assert_eq!(usage.objects_total_count, 2, "{kind}");
assert_eq!(usage.buckets_usage[&cold].objects_count, 1, "{kind}");
}
clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_entry_fallback_covers_overflow_and_new_bucket_inventory() {
let (_dir, store) = setup_two_pool_scanner_store().await;
clear_dirty_usage_buckets_for_tests();
let hot = format!("hot-{}", Uuid::new_v4().simple());
create_bucket(&store, &hot).await;
record_dirty_usage_bucket(&hot);
let baseline = run_entry(&store, 1, None, true).await;
persist_baseline(&store, &baseline).await;
for index in 0..=crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES {
record_dirty_usage_bucket(&format!("overflow-{index}"));
}
assert!(dirty_usage_buckets_for_tests().len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES);
let usage = run_entry(&store, 2, None, true).await;
assert_eq!(usage.objects_total_count, 1);
clear_dirty_usage_buckets_for_tests();
record_dirty_usage_bucket(&hot);
let new_bucket = format!("new-{}", Uuid::new_v4().simple());
create_bucket(&store, &new_bucket).await;
// Even a previously valid baseline cannot cover the changed inventory.
let usage = run_entry(&store, 3, None, true).await;
assert_eq!(usage.objects_total_count, 2);
assert_eq!(usage.buckets_usage[&new_bucket].objects_count, 1);
clear_dirty_usage_buckets_for_tests();
}
+1 -1
View File
@@ -17,7 +17,7 @@ A multi-pool layout in which any pool spans several localhost ports is not expre
Data-movement cases fail closed. A decommission or rebalance test must observe a successful start response, an active state, a clean terminal state, non-zero movement counters, and post-operation object integrity. An unsupported response, HTTP 5xx, missing status fields, cleanup warning, or zero-progress terminal response fails the case; pre/post S3 availability alone is not evidence that movement ran.
The four expansion pools must report independent capacity. Four directories on one runner filesystem all return the same `statfs` totals, so RustFS correctly concludes that no pool is less free than the cluster average and performs no rebalance. The Actions job runs on GitHub-hosted `ubuntu-latest` and mounts four isolated 1 GiB tmpfs filesystems, then exports their absolute paths through `RUSTFS_E2E_POOL_ROOTS`. It does not use the self-hosted `sm-standard-4` ARC pods: those cannot create filesystems (`mount -o loop` fails with `No such file or directory`, and `mount -t tmpfs` fails with `cannot mount tmpfs read-only`). Sized tmpfs still reports a distinct `st_dev` and independent 1 GiB `statfs` capacity. The harness rejects missing, duplicate, relative, nonexistent, or same-device roots instead of allowing a vacuous movement pass. Planned pool additions stop every process with SIGTERM; hard process termination remains a chaos-only fault. After the fourth pool joins, the harness performs one full graceful persistent restart: this proves the expanded pool map survives restart and ensures movement begins only after every replica can load the converged metadata.
The four expansion pools must report independent capacity. Four directories on one runner filesystem all return the same `statfs` totals, so RustFS correctly concludes that no pool is less free than the cluster average and performs no rebalance. The Actions job mounts four isolated 1 GiB tmpfs filesystems and exports their absolute paths through `RUSTFS_E2E_POOL_ROOTS`. It does not use ext4 loop devices: the `sm-standard-4` ARC pods have no `/dev/loop-control`, so `mount -o loop` fails with `No such file or directory`. Sized tmpfs still reports a distinct `st_dev` and independent 1 GiB `statfs` capacity. The harness rejects missing, duplicate, relative, nonexistent, or same-device roots instead of allowing a vacuous movement pass. Planned pool additions stop every process with SIGTERM; hard process termination remains a chaos-only fault. After the fourth pool joins, the harness performs one full graceful persistent restart: this proves the expanded pool map survives restart and ensures movement begins only after every replica can load the converged metadata.
The expansion fixture is an all-current-binary fleet, so it initializes pool metadata with the documented V3 write and fleet-confirmation gates. Decommission cases write their baseline objects, version history, and multipart data into pool 0 before adding pools 13, then retire pool 0. This makes a passing result evidence of user-data movement rather than merely an internal-metadata counter changing.
+2 -250
View File
@@ -44,9 +44,8 @@ use super::storage_api::HTTPRangeSpec;
use super::storage_api::remote_s3_client::RemoteS3ClientError;
use hmac::{Hmac, Mac, digest::KeyInit};
use http::{HeaderMap, HeaderValue, Method};
use percent_encoding::percent_decode_str;
use quick_xml::Reader;
use quick_xml::events::{BytesStart, Event};
use quick_xml::events::Event;
use sha2::Sha256;
use std::collections::{BTreeMap, HashMap};
use url::Url;
@@ -380,23 +379,13 @@ fn parse_list_blobs(xml: &str) -> Result<AzureListing, SourceError> {
_ => {
let end = start.to_end().into_owned();
let text = leaf_text(&mut reader, end.name())?;
let text = if name == "name" {
decode_list_name(&start, text)?
} else {
text
};
apply_list_field(&name, text, &mut blob, &mut prefixes, &mut next_marker, in_blob_prefix);
}
}
}
Ok(Event::Empty(empty)) => {
let name = local_name(empty.name().as_ref());
let text = if name == "name" {
decode_list_name(&empty, String::new())?
} else {
String::new()
};
apply_list_field(&name, text, &mut blob, &mut prefixes, &mut next_marker, in_blob_prefix);
apply_list_field(&name, String::new(), &mut blob, &mut prefixes, &mut next_marker, in_blob_prefix);
}
Ok(Event::End(end)) => match local_name(end.name().as_ref()).as_str() {
"blob" => {
@@ -437,43 +426,6 @@ fn parse_list_blobs(xml: &str) -> Result<AzureListing, SourceError> {
})
}
/// Azure marks XML-inexpressible blob/prefix names with `Encoded="true"`.
/// Only those names are URI-decoded, once; ordinary percent signs and `+`
/// are part of the key, and NextMarker remains an opaque cursor.
fn decode_list_name(start: &BytesStart<'_>, text: String) -> Result<String, SourceError> {
let mut encoded = false;
for attribute in start.attributes() {
let attribute = attribute.map_err(|_| SourceError::Other("source listing name has invalid attributes".to_string()))?;
if attribute.key.as_ref() == "Encoded" {
let value = quick_xml::escape::unescape(&attribute.value)
.map_err(|_| SourceError::Other("source listing name has an invalid Encoded attribute".to_string()))?;
encoded = match value.as_ref() {
"true" | "1" => true,
"false" | "0" => false,
_ => return Err(SourceError::Other("source listing name has an invalid Encoded attribute".to_string())),
};
}
}
if !encoded {
return Ok(text);
}
// percent_decode_str leaves malformed escapes untouched. Refuse them
// rather than return a different key or replace invalid UTF-8 with U+FFFD.
let mut bytes = text.bytes();
while let Some(byte) = bytes.next() {
if byte == b'%'
&& !(bytes.next().is_some_and(|b| b.is_ascii_hexdigit()) && bytes.next().is_some_and(|b| b.is_ascii_hexdigit()))
{
return Err(SourceError::Other("source listing name has invalid percent encoding".to_string()));
}
}
percent_decode_str(&text)
.decode_utf8()
.map(|name| name.into_owned())
.map_err(|_| SourceError::Other("source listing name is not valid UTF-8".to_string()))
}
fn apply_list_field(
name: &str,
text: String,
@@ -634,13 +586,6 @@ mod tests {
const LAST_PAGE: &str = r#"<?xml version="1.0" encoding="utf-8"?>
<EnumerationResults><Blobs><Blob><Name>only.txt</Name><Properties><Content-Length>1</Content-Length></Properties></Blob></Blobs><NextMarker /></EnumerationResults>"#;
const ENCODED_NAME_PAGE: &str = r#"<EnumerationResults><Blobs>
<Blob><Name Encoded="true">%EF%BF%BE/part%252F+%20%26.txt</Name><Properties><Content-Length>5</Content-Length></Properties></Blob>
<Blob><Name>%EF%BF%BE/part%252F+%20%26.txt</Name><Properties><Content-Length>5</Content-Length></Properties></Blob>
<BlobPrefix><Name Encoded="true">%EF%BF%BF%2F</Name></BlobPrefix>
<BlobPrefix><Name Encoded="false">literal%FF+/</Name></BlobPrefix>
</Blobs><NextMarker Encoded="true">opaque%2F+cursor</NextMarker></EnumerationResults>"#;
const TAGS: &str = r#"<?xml version="1.0" encoding="utf-8"?>
<Tags><TagSet>
<Tag><Key>env</Key><Value>prod</Value></Tag>
@@ -673,87 +618,6 @@ mod tests {
let listing = parse_list_blobs(LAST_PAGE).expect("page should parse");
assert_eq!(listing.objects.len(), 1);
assert!(listing.next_marker.is_none(), "an empty NextMarker is not a cursor");
let empty = parse_list_blobs("<EnumerationResults><Blobs/><NextMarker/></EnumerationResults>")
.expect("an empty final page is valid");
assert!(empty.objects.is_empty());
assert!(empty.prefixes.is_empty());
assert!(empty.next_marker.is_none());
}
#[test]
fn list_blobs_decodes_only_marked_names_once() {
let listing = parse_list_blobs(ENCODED_NAME_PAGE).expect("encoded names should parse");
assert_eq!(listing.objects[0].key, "\u{fffe}/part%2F+ &.txt");
assert_eq!(listing.objects[1].key, "%EF%BF%BE/part%252F+%20%26.txt");
assert_eq!(listing.prefixes, ["\u{ffff}/", "literal%FF+/"]);
assert_eq!(listing.next_marker.as_deref(), Some("opaque%2F+cursor"));
for (attribute, text, expected) in [
("", "a%2Fb+ &amp;.txt", "a%2Fb+ &.txt"),
("Encoded=\"false\"", "a%2Fb+ &amp;.txt", "a%2Fb+ &.txt"),
("Encoded=\"0\"", "a%2Fb+ &amp;.txt", "a%2Fb+ &.txt"),
("Encoded=\"true\"", "a%2Fb+ &amp;.txt", "a/b+ &.txt"),
("Encoded=\"1\"", "a%2Fb+ &amp;.txt", "a/b+ &.txt"),
("Encoded=\"tr&#117;e\"", "a%2Fb+ &amp;.txt", "a/b+ &.txt"),
("", "%", "%"),
("Encoded=\"false\"", "%", "%"),
("", "中文/plain%2F+name%", "中文/plain%2F+name%"),
("Encoded=\"false\"", "中文/plain%2F+name%", "中文/plain%2F+name%"),
(
"Encoded=\"true\"",
"%EF%BF%BE%EF%BF%BF/%E4%B8%AD%E6%96%87-%25-%2B-%252F+&amp;-%26amp%3B",
"\u{fffe}\u{ffff}/中文-%-+-%2F+&-&amp;",
),
] {
for container in ["Blob", "BlobPrefix"] {
let properties = if container == "Blob" {
"<Properties><Content-Length>0</Content-Length></Properties>"
} else {
""
};
let xml = format!(
"<EnumerationResults><Blobs><{container}><Name {attribute}>{text}</Name>{properties}</{container}></Blobs></EnumerationResults>"
);
let listing = parse_list_blobs(&xml).expect("valid name");
if container == "Blob" {
assert_eq!(listing.objects.len(), 1);
assert_eq!(listing.objects[0].key, expected, "{container}: {attribute}, {text}");
assert_eq!(listing.objects[0].size, 0, "a named zero-byte blob remains valid");
} else {
assert!(listing.objects.is_empty(), "a prefix-only page remains valid");
assert_eq!(listing.prefixes, [expected], "{container}: {attribute}, {text}");
}
}
}
}
#[test]
fn list_blobs_rejects_invalid_encoded_names_without_returning_partial_entries() {
for name in [
"<Name Encoded=\"true\">%</Name>",
"<Name Encoded=\"true\">%2</Name>",
"<Name Encoded=\"true\">%GG</Name>",
"<Name Encoded=\"true\">%FF</Name>",
"<Name Encoded=\"true\">%E2%82</Name>",
"<Name Encoded=\"true\">%C0%AF</Name>",
"<Name Encoded=\"true\">%ED%A0%80</Name>",
"<Name Encoded=\"maybe\">a</Name>",
"<Name Encoded=\"true\" Encoded=\"false\">a</Name>",
"<Name Encoded=\"true\" Encoded=\"false\"/>",
"<Name Encoded=\"&unknown;\">a</Name>",
] {
for container in ["Blob", "BlobPrefix"] {
let properties = if container == "Blob" {
"<Properties><Content-Length>1</Content-Length></Properties>"
} else {
""
};
let xml = format!(
"<EnumerationResults><Blobs><Blob><Name>before</Name><Properties><Content-Length>1</Content-Length></Properties></Blob><{container}>{name}{properties}</{container}><Blob><Name>after</Name><Properties><Content-Length>1</Content-Length></Properties></Blob></Blobs><NextMarker>opaque%2B+marker</NextMarker></EnumerationResults>"
);
assert!(matches!(parse_list_blobs(&xml), Err(SourceError::Other(_))), "{container}: {name}");
}
}
}
#[test]
@@ -1042,118 +906,6 @@ mod tests {
);
}
#[tokio::test]
async fn listed_encoded_and_literal_names_get_distinct_source_objects() {
let mut head_headers = blob_headers();
head_headers.push(("Content-Length", "5".to_string()));
let (endpoint, recorded) = scripted_server(vec![
ScriptedResponse::new(200, Vec::new(), ENCODED_NAME_PAGE.to_string()),
ScriptedResponse::new(200, head_headers.clone(), String::new()),
ScriptedResponse::new(200, blob_headers(), "first".to_string()),
ScriptedResponse::new(200, head_headers, String::new()),
ScriptedResponse::new(200, blob_headers(), "other".to_string()),
])
.await;
let backend = backend(&endpoint, Credential::SharedKey(vec![7_u8; 32]));
let page = backend
.list(&SourceListRequest {
max_keys: 4,
..Default::default()
})
.await
.expect("list names");
assert_eq!(page.objects.len(), 2);
assert_eq!(page.objects[0].key, "\u{fffe}/part%2F+ &.txt");
assert_eq!(page.objects[1].key, "%EF%BF%BE/part%252F+%20%26.txt");
assert_eq!(page.common_prefixes, ["\u{ffff}/", "literal%FF+/"]);
for (object, body) in page.objects.iter().zip([b"first", b"other"]) {
let head = backend.head(&object.key).await.expect("head listed object");
assert_eq!(head.size, object.size);
let got = backend.get(&object.key, None).await.expect("get listed object");
assert_eq!(got.body.collect().await.expect("source body").into_bytes().as_ref(), body);
}
let recorded = recorded.lock().expect("recorder lock");
assert_eq!(recorded.len(), 5);
for (requests, expected) in recorded[1..].chunks_exact(2).zip([
"/legacy/%EF%BF%BE/part%252F+%20&.txt",
"/legacy/%25EF%25BF%25BE/part%25252F+%2520%2526.txt",
]) {
assert_eq!(requests[0].method, "HEAD");
assert_eq!(requests[1].method, "GET");
assert_eq!(requests[0].target, expected);
assert_eq!(requests[1].target, expected);
}
}
#[tokio::test]
async fn listed_encoded_prefix_and_opaque_marker_round_trip_through_query_encoding() {
const PAGE: &str = r#"<EnumerationResults><Blobs><BlobPrefix><Name Encoded="true">%EF%BF%BE%EF%BF%BF/%E4%B8%AD%E6%96%87/%252F%2B%25+%26/</Name></BlobPrefix></Blobs><NextMarker>opaque%2B+marker</NextMarker></EnumerationResults>"#;
let (endpoint, recorded) = scripted_server(vec![
ScriptedResponse::new(200, Vec::new(), PAGE.to_string()),
ScriptedResponse::new(200, Vec::new(), LAST_PAGE.to_string()),
ScriptedResponse::new(200, Vec::new(), LAST_PAGE.to_string()),
])
.await;
let backend = backend(&endpoint, Credential::SharedKey(vec![7_u8; 32]));
let first = backend
.list(&SourceListRequest {
delimiter: Some("/"),
max_keys: 1,
..Default::default()
})
.await
.expect("list encoded prefix");
assert!(first.objects.is_empty());
assert_eq!(first.common_prefixes, ["\u{fffe}\u{ffff}/中文/%2F+%+&/"]);
assert_eq!(first.next_continuation_token.as_deref(), Some("opaque%2B+marker"));
assert!(first.is_truncated);
let second = backend
.list(&SourceListRequest {
delimiter: Some("/"),
continuation_token: first.next_continuation_token.as_deref(),
max_keys: 1,
..Default::default()
})
.await
.expect("continue with the original listing conditions");
assert!(!second.is_truncated);
assert!(second.next_continuation_token.is_none());
let nested = backend
.list(&SourceListRequest {
prefix: Some(&first.common_prefixes[0]),
delimiter: Some("/"),
max_keys: 1,
..Default::default()
})
.await
.expect("start a separate listing under the returned logical prefix");
assert!(!nested.is_truncated);
let recorded = recorded.lock().expect("recorder lock");
assert_eq!(recorded.len(), 3);
assert_eq!(recorded[0].method, "GET");
assert!(!recorded[0].target.contains("marker="));
assert_eq!(recorded[1].method, "GET");
assert_eq!(
recorded[1].target,
"/legacy?restype=container&comp=list&delimiter=%2F&marker=opaque%252B%2Bmarker&maxresults=1"
);
let request_url = endpoint.join(&recorded[1].target).expect("recorded request URL");
let query: HashMap<_, _> = request_url.query_pairs().into_owned().collect();
assert_eq!(query.get("marker").map(String::as_str), Some("opaque%2B+marker"));
assert_eq!(recorded[2].method, "GET");
assert_eq!(
recorded[2].target,
"/legacy?restype=container&comp=list&prefix=%EF%BF%BE%EF%BF%BF%2F%E4%B8%AD%E6%96%87%2F%252F%2B%25%2B%26%2F&delimiter=%2F&maxresults=1"
);
let request_url = endpoint.join(&recorded[2].target).expect("recorded prefix request URL");
let query: HashMap<_, _> = request_url.query_pairs().into_owned().collect();
assert_eq!(query.get("prefix"), Some(&first.common_prefixes[0]));
assert!(!query.contains_key("marker"));
}
#[tokio::test]
async fn tagging_and_probe_address_the_right_resources() {
let (endpoint, recorded) = scripted_server(vec![