mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 03:59:14 +00:00
Compare commits
32 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d603ea3922 | |||
| 4068f8440c | |||
| 982ed2888f | |||
| eec6d3a76d | |||
| bbca907781 | |||
| 51982be687 | |||
| b8f0119320 | |||
| 20ca988e85 | |||
| 045770d09d | |||
| 278cbaefd5 | |||
| 982bebf6a5 | |||
| 1aafe477ec | |||
| c33a42c5ab | |||
| e3fd1ae19a | |||
| 6c1d401215 | |||
| c80ff445c0 | |||
| 3d4c05ddf9 | |||
| f9dd0388fa | |||
| 527860d71e | |||
| 7d2c120073 | |||
| 489c3276a4 | |||
| b09a4f3706 | |||
| a0c32e7c6b | |||
| d5a1530409 | |||
| e7b0349a4f | |||
| bdbdca07c8 | |||
| 53efaa2b8f | |||
| ec9672a397 | |||
| cee35f7e54 | |||
| ef7e7afd8c | |||
| 652ebb12c6 | |||
| 38c03d9d5d |
@@ -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
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
@@ -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 1–3, then retire pool 0. This makes a passing result evidence of user-data movement rather than merely an internal-metadata counter changing.
|
||||
|
||||
|
||||
@@ -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+ &.txt", "a%2Fb+ &.txt"),
|
||||
("Encoded=\"false\"", "a%2Fb+ &.txt", "a%2Fb+ &.txt"),
|
||||
("Encoded=\"0\"", "a%2Fb+ &.txt", "a%2Fb+ &.txt"),
|
||||
("Encoded=\"true\"", "a%2Fb+ &.txt", "a/b+ &.txt"),
|
||||
("Encoded=\"1\"", "a%2Fb+ &.txt", "a/b+ &.txt"),
|
||||
("Encoded=\"true\"", "a%2Fb+ &.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+&-%26amp%3B",
|
||||
"\u{fffe}\u{ffff}/中文-%-+-%2F+&-&",
|
||||
),
|
||||
] {
|
||||
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![
|
||||
|
||||
Reference in New Issue
Block a user