chore(ci): integrate current main for E2E provenance

This commit is contained in:
overtrue
2026-09-08 23:55:05 +08:00
20 changed files with 1693 additions and 37 deletions
@@ -82,6 +82,12 @@ jobs:
performance-test:
runs-on: pf-testing
timeout-minutes: 900
env:
RUSTFS_BENCH_SCRIPT: ${{ github.workspace }}/auto-testing/rustfs_performance_testing.sh
RUSTFS_WARP_METHODS: ${{ inputs.test_method }}
RUSTFS_WARP_SIZES: ${{ inputs.object_size }}
RUSTFS_WARP_DURATION: ${{ inputs.warp_duration || '5m' }}
RUSTFS_WARP_CONCURRENCY: ${{ inputs.warp_concurrency || '64' }}
# Run on manual dispatch, or when the nightly build completed successfully.
# Skipped when nightly failed.
if: ${{ github.event_name == 'workflow_dispatch' || github.event_name == 'repository_dispatch' }}
@@ -158,19 +164,15 @@ jobs:
- name: Run benchmark (GET/PUT/MIXED)
id: benchmark
run: |
# Empty on automatic (workflow_run) runs -> full 30 rounds.
# Manual dispatch can restrict method(s)/size(s).
export WARP_METHODS="${{ inputs.test_method }}"
export WARP_SIZES="${{ inputs.object_size }}"
./auto-testing/rustfs_performance_test.sh \
--step 5 -y \
--warp-duration "${{ inputs.warp_duration || '5m' }}" \
--warp-concurrency "${{ inputs.warp_concurrency || '64' }}" \
--log-file "${LOG_FILE}"
- name: Analyze results
if: ${{ steps.benchmark.conclusion == 'success' }}
run: |
export WARP_METHODS="${RUSTFS_WARP_METHODS}" WARP_SIZES="${RUSTFS_WARP_SIZES}"
export WARP_DURATION="${RUSTFS_WARP_DURATION}" WARP_CONCURRENCY="${RUSTFS_WARP_CONCURRENCY}"
./auto-testing/rustfs_performance_test.sh --step 6 -y --log-file "${LOG_FILE:-/dev/null}"
- name: Collect RustFS version info
+1
View File
@@ -52,6 +52,7 @@ docs
__pycache__/
!docs/
docs/*
!docs/README.md
!docs/architecture/
!docs/architecture/**
!docs/operations/
+10
View File
@@ -7,6 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
### Replication
- Object Lock replication PUTs now carry a required integrity header, fixing target rejection introduced by the plain-payload default ([#7097](https://github.com/rustfs/rustfs/pull/7097)). This changes the default outbound request for locked objects but adds no persisted format.
- Multipart source objects stay on the multipart transport even when their checksum record is a whole-object checksum, so objects above the single-PUT limit remain replicable ([#7047](https://github.com/rustfs/rustfs/pull/7047)).
- Targets that mint their own version IDs now use a per-target version ledger for tag, retention, legal-hold, and permanent-delete mutations; ambiguous pre-ledger matches fail with backoff instead of guessing ([#7368](https://github.com/rustfs/rustfs/pull/7368)). This adds dual-prefixed internal metadata keys that older readers ignore.
- Single-part source checksums are forwarded as `x-amz-checksum-*` headers instead of user metadata, so the replica preserves checksum responses ([#7313](https://github.com/rustfs/rustfs/pull/7313)). This changes the default outbound headers for checksummed objects.
- Site-replication outage recovery now uses a bounded 30-second retry drain plus the 600-second full reconciliation pass, persists destructive liabilities before local deletion, and fences replay settlement and peer edits ([#7148](https://github.com/rustfs/rustfs/pull/7148)). Persisted additions are optional and ignored by older readers.
- IAM snapshot/deletion replay, target-assigned delete-marker purges, timestamp ordering, and best-effort peer broadcast now close the control-plane gaps found by the R6 review ([#7195](https://github.com/rustfs/rustfs/pull/7195)).
- Upgrade and rollback: upgrade every node in one site consecutively and verify reconciliation before moving to the next site; do not intentionally run a site mixed-version. Target-version ledger keys are harmless on rollback, although old code cannot use their routing. Before rolling back past [#7307](https://github.com/rustfs/rustfs/pull/7307), drain or repair every pending version purge: older code can free a retained version's data directory before its remote purge is acknowledged. See `docs/operations/site-replication-operations.md`.
### Security
- **Presigned URLs honour only signed headers** (GHSA-g8w9-qw9q-fghr): a SigV4 presigned request that carries an `x-amz-*` request header not listed in `X-Amz-SignedHeaders` is now rejected with `403 AccessDenied` ("There were headers present in the request which were not signed"), matching AWS S3. Previously the holder of a presigned `PutObject` URL could add unsigned `x-amz-tagging`, `x-amz-storage-class`, `x-amz-website-redirect-location`, ACL, metadata, Object Lock or SSE headers and have them applied. Presigners that intend a property must set it before signing so the SDK lists the header in `SignedHeaders`; `x-amz-cf-id` (CloudFront) remains tolerated unsigned. Header-signed SigV4 and SigV2 requests are unchanged.
@@ -4235,6 +4235,16 @@ async fn test_bucket_replication_acceptance_matrix_local_dual_targets() -> TestR
<ExistingObjectReplication><Status>Enabled</Status></ExistingObjectReplication>
<Destination><Bucket>{target_b_arn}</Bucket></Destination>
</Rule>
<Rule>
<ID>matrix-and-tags</ID>
<Priority>135</Priority>
<Status>Enabled</Status>
<Filter><And><Prefix>and-tags/</Prefix><Tag><Key>env</Key><Value>prod</Value></Tag><Tag><Key>tier</Key><Value>gold</Value></Tag></And></Filter>
<DeleteMarkerReplication><Status>Disabled</Status></DeleteMarkerReplication>
<DeleteReplication><Status>Enabled</Status></DeleteReplication>
<ExistingObjectReplication><Status>Enabled</Status></ExistingObjectReplication>
<Destination><Bucket>{target_b_arn}</Bucket></Destination>
</Rule>
<Rule>
<ID>matrix-disabled</ID>
<Priority>140</Priority>
@@ -4289,6 +4299,7 @@ async fn test_bucket_replication_acceptance_matrix_local_dual_targets() -> TestR
"matrix-prefix",
"matrix-tag",
"matrix-disabled",
"matrix-and-tags",
"matrix-priority-high",
"Priority>200",
"<Status>Disabled</Status>",
@@ -4409,6 +4420,30 @@ async fn test_bucket_replication_acceptance_matrix_local_dual_targets() -> TestR
put_single_tag_current(&source_client, source_bucket, "tagged/no-match.txt", "route", "tagged").await?;
assert_replication_key_absent(&target_client_b, target_bucket_b, "tagged/no-match.txt", Duration::from_secs(3)).await?;
// S3 and MinIO both read `And.Tags` as AND: an object carrying only one of
// the required tags is not admitted. Matching any single tag would push
// data to a destination the rule never selected (backlog#2366 P1-1), and
// the two-tag rule is the shape `mc replicate add --tags "k1=v1&k2=v2"`
// writes, so a single-tag rule passing is not evidence for this.
source_client
.put_object()
.bucket(source_bucket)
.key("and-tags/partial.txt")
.tagging("env=prod")
.body(ByteStream::from_static(b"one of two tags"))
.send()
.await?;
assert_replication_key_absent(&target_client_b, target_bucket_b, "and-tags/partial.txt", Duration::from_secs(3)).await?;
source_client
.put_object()
.bucket(source_bucket)
.key("and-tags/full.txt")
.tagging("env=prod&tier=gold")
.body(ByteStream::from_static(b"both tags"))
.send()
.await?;
wait_for_user_get_object(&target_client_b, target_bucket_b, "and-tags/full.txt").await?;
source_client
.put_object()
.bucket(source_bucket)
+60 -1
View File
@@ -166,6 +166,24 @@ fn rule_replicates(rule: &ReplicationRule, obj: &ObjectOpts) -> bool {
}
}
fn replication_filter_tags_match(filter: &s3s::dto::ReplicationRuleFilter, object_tags: &HashMap<String, String>) -> bool {
let tag_matches = |tag: &s3s::dto::Tag| match (&tag.key, &tag.value) {
(None, None) => true,
(Some(key), _) if key.is_empty() => true,
(Some(key), Some(value)) => object_tags.get(key) == Some(value),
_ => false,
};
filter
.and
.as_ref()
.and_then(|and| and.tags.as_deref())
.into_iter()
.flatten()
.chain(filter.tag.iter())
.all(tag_matches)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReplicationTargetValidationError {
RoleWithMultipleDestinations,
@@ -704,7 +722,7 @@ impl ReplicationConfigurationExt for ReplicationConfiguration {
if let Some(filter) = &rule.filter {
let object_tags = ReplicationTagFilter::decode_tags_to_map(&obj.user_tags);
if filter.test_tags(&object_tags) {
if replication_filter_tags_match(filter, &object_tags) {
rules.push(rule.clone());
}
} else {
@@ -1139,6 +1157,47 @@ mod tests {
assert_eq!(validate_replication_config_structure(&structure_config(vec![rule])), Ok(()));
}
#[test]
fn actionable_rules_require_every_and_tag_to_match() {
let mut rule = replication_rule("rule-1", "arn:target:a");
rule.filter = Some(s3s::dto::ReplicationRuleFilter {
and: Some(s3s::dto::ReplicationRuleAndOperator {
prefix: None,
tags: Some(vec![
s3s::dto::Tag {
key: Some("env".to_string()),
value: Some("prod".to_string()),
},
s3s::dto::Tag {
key: Some("tier".to_string()),
value: Some("gold".to_string()),
},
]),
}),
..Default::default()
});
let config = structure_config(vec![rule]);
let object = |user_tags: &str| ObjectOpts {
name: "object".to_string(),
user_tags: user_tags.to_string(),
..Default::default()
};
assert!(config.filter_target_arns(&object("env=prod")).is_empty());
assert_eq!(config.filter_target_arns(&object("env=prod&tier=gold")), vec!["arn:target:a"]);
assert!(config.filter_target_arns(&object("")).is_empty());
let mut malformed = config;
malformed.rules[0].filter.as_mut().unwrap().and.as_mut().unwrap().tags = Some(vec![s3s::dto::Tag {
key: Some("env".to_string()),
value: None,
}]);
assert!(
malformed.filter_target_arns(&object("env=prod")).is_empty(),
"a malformed tag filter must fail closed"
);
}
#[test]
fn structure_validation_allows_tag_filter_when_delete_marker_replication_disabled() {
let mut rule = replication_rule("rule-1", "arn:target:a");
+44
View File
@@ -580,6 +580,30 @@ impl FailStats {
FailedMetric { count, size }
}
/// Both rolling windows from one walk of the samples. `short` must be the
/// narrower window; the walk stops at `long`. Callers that need both (the
/// per-node site snapshot) would otherwise scan the deque twice while
/// holding the bucket-stats read lock, and the deque is only bounded by
/// the one-hour window - an unreachable target under load fills it.
pub fn recent_windows(&self, short: Duration, long: Duration) -> (FailedMetric, FailedMetric) {
let now = Instant::now();
let mut short_metric = FailedMetric::default();
let mut long_metric = FailedMetric::default();
for sample in self.recent.iter().rev() {
let age = now.duration_since(sample.observed_at);
if age > long {
break;
}
if age <= short {
short_metric.count += 1;
short_metric.size += sample.size;
}
long_metric.count += 1;
long_metric.size += sample.size;
}
(short_metric, long_metric)
}
pub fn merge(&self, other: &FailStats) -> Self {
Self {
count: self.count.saturating_add(other.count),
@@ -912,6 +936,26 @@ mod tests {
assert_eq!(last_hour.size, 96);
}
#[test]
fn fail_stats_recent_windows_matches_two_separate_scans() {
let mut stats = FailStats::default();
stats.add_size(64, None::<&()>);
stats.add_size(32, None::<&()>);
let (minute, hour) = stats.recent_windows(Duration::from_secs(60), Duration::from_secs(60 * 60));
let expected_minute = stats.recent_since(Duration::from_secs(60));
let expected_hour = stats.recent_since(Duration::from_secs(60 * 60));
assert_eq!((minute.count, minute.size), (expected_minute.count, expected_minute.size));
assert_eq!((hour.count, hour.size), (expected_hour.count, expected_hour.size));
assert_eq!(minute.count, 2);
assert_eq!(hour.size, 96);
let empty = FailStats::default();
let (minute, hour) = empty.recent_windows(Duration::from_secs(60), Duration::from_secs(60 * 60));
assert_eq!((minute.count, minute.size, hour.count, hour.size), (0, 0, 0, 0));
}
#[test]
fn fail_stats_saturate_instead_of_wrapping() {
let mut stats = FailStats {
+23
View File
@@ -0,0 +1,23 @@
# Documentation
Use the focused indexes rather than treating this directory as an unordered
collection:
- [Architecture knowledge base](architecture/README.md)
- [Testing references](testing/README.md)
## Operations
Operational runbooks live under [`operations/`](operations/). Replication
operators should start with:
| Runbook | Use it for |
|---|---|
| [Site replication operations](operations/site-replication-operations.md) | Health fields, pending operations, outage recovery, re-pair admission, IAM/SSE boundaries, and upgrades. |
| [Replication target check](operations/replication-check.md) | Validating an S3 destination and version fidelity before enabling replication. |
| [Replication object size limits](operations/replication-object-size-limits.md) | Multipart routing, large-object limits, and retry characteristics. |
| [Replication outbound transport](operations/replication-outbound-transport.md) | Integrity headers, generic target behavior, and transport knobs. |
Other runbooks remain grouped by filename in [`operations/`](operations/);
architecture pages link to the relevant runbook where a cross-boundary
procedure is required.
+3 -1
View File
@@ -60,6 +60,8 @@ Required headings and strings in these files are asserted by `scripts/check_arch
| [minio-rustfs-router-compatibility.md](minio-rustfs-router-compatibility.md) | a client or `mc` call that works against MinIO fails against RustFS and you need to know whether the endpoint is missing, stubbed, or deliberately different |
| [minio-file-format-compat.md](minio-file-format-compat.md) | deciding whether a MinIO drive set, bucket-metadata blob, or SSE object can be read or imported by a given RustFS build, or before touching a listed version anchor |
Operations runbooks live in [../operations/](../operations/) and testing references in [../testing/README.md](../testing/README.md).
Operations runbooks are registered in the [documentation operations index](../README.md#operations), and testing references live in [../testing/README.md](../testing/README.md).
For replication operations, start with [site replication operations](../operations/site-replication-operations.md), [replication target check](../operations/replication-check.md), [replication object size limits](../operations/replication-object-size-limits.md), and [replication outbound transport](../operations/replication-outbound-transport.md).
For per-node HTTP failure ratios and cached storage probe provenance, see [S3 write failure diagnostics](../operations/s3-write-failure-diagnostics.md).
@@ -38,6 +38,42 @@ Counts ignore blank lines and comments; compute them from the files. The lifecyc
"Supported" for the SSE row means RustFS encrypts and decrypts its own objects. MinIO SSE objects (SSE-S3, SSE-KMS, SSE-C) are not readable in default builds; see [minio-file-format-compat.md Part C](minio-file-format-compat.md#part-c--server-side-encryption-sse) for the `rio-v2` migration build.
## Replication Support Boundary
Site replication and bucket replication are not the same compatibility claim.
Site replication requires RustFS-compatible peer admin APIs and coordinates
IAM, topology, buckets, and metadata. A generic S3-compatible service can only
be a bucket-replication data target.
For a generic S3 target, RustFS supports object PUT/HEAD/DELETE, multipart
uploads, tags, version deletes, and Object Lock mutations when the target
implements the corresponding S3 APIs and has versioning enabled. Targets that
mint their own version IDs are supported through a per-target version ledger;
pre-ledger replicas are adopted only when exact key and ETag identify one
unambiguous target version. `NoSuchVersion` for an already absent addressed
replica is treated as converged.
The following are capability boundaries, not universal S3 claims:
- `GET /BUCKET?replication-check` must pass the phases required by the intended
workload. `VersionFidelity` may report a minting target as mismatched even
though ledger-addressed delete and Object Lock phases succeed.
- A target that rejects standard multipart constraints, required Object Lock
integrity headers, or the configured checksum framing is unsupported until
its transport settings are made compatible.
- SSE-S3 and SSE-KMS are decrypted at the source and re-encrypted by the
destination's KMS. SSE-C uses ciphertext passthrough and requires target
evidence. Unsupported or ambiguous encryption metadata fails closed.
- ACL authorization is intentionally unsupported, and generic targets never
receive RustFS IAM/site-control-plane state.
- RustFS does not guess between multiple target versions with the same key and
ETag. The mutation remains failed and retryable until repair establishes an
unambiguous mapping.
See [site replication operations](../operations/site-replication-operations.md)
for health, recovery, and upgrade rules and [replication outbound transport](../operations/replication-outbound-transport.md)
for the tested target classes and knobs.
## Not Yet Passing
Standard S3 areas that must not be described as complete:
@@ -0,0 +1,258 @@
# Site Replication Operations
**Use this when:** operating a site-replication deployment, diagnosing a peer
outage or incomplete topology change, pairing sites that already contain data,
or planning an upgrade.
**Source of truth:** `rustfs/src/admin/handlers/site_replication.rs`,
`rustfs/src/site_replication/`, and the bucket-replication worker under
`crates/ecstore/src/bucket/replication/`.
Site replication combines two different convergence paths:
- the control plane replicates buckets, bucket metadata, IAM, and topology;
- ordinary bucket replication moves object versions and delete operations.
An `enabled: true` response only says that a site has more than one configured
peer. It does not prove that every peer is reachable or caught up. Always read
`pendingOperation`, `retryStats`, `PeerErrors`, and `Metrics` as well.
## Routine checks
Run these commands from an admin workstation with one alias per site:
```console
mc admin replicate info site-a
mc admin replicate status site-a
```
Check more than one site. A partition can leave each side with a different but
locally valid view.
`replicate info` is the compact control-plane view:
| Field | Interpretation |
|---|---|
| `enabled` | More than one site is configured; this is not a health verdict. |
| `sites` | The locally persisted topology. Compare deployment IDs and endpoints on every site. |
| `retryStats.pending` | Collapsed peer deliveries waiting to be retried. |
| `retryStats.failed` | Deliveries that crossed the escalation threshold and require attention. |
| `retryStats.lastError` | A redacted summary of the most recent delivery failure. |
| `pendingOperation` | A durable multi-step topology operation described below. Absence is the healthy steady state. |
`replicate status` adds detailed convergence state:
| Field | Interpretation |
|---|---|
| `Sites` / `PeerStates` | Configured peers and derived reachability/configuration state. |
| `PeerErrors` | A peer could not be queried. Its detailed counters may be absent; do not read zeros as success. |
| `BucketStats` | Per-bucket presence and versioning, replication, lifecycle, Object Lock, and metadata mismatches. |
| `PolicyStats`, `UserStats`, `GroupStats` | IAM inventory mismatches. |
| `RetryStats` | Durable control-plane retry backlog and escalation count. |
| `Metrics.replMetrics` | Per-destination online state, downtime, replicated counts/bytes, and `failed` totals/windows. |
| `Metrics.queued` / `Metrics.inProgress` | Object work waiting or active on the responding node. |
| `Metrics.errors` | Node-level object-replication failures. When only queue statistics are available, RustFS synthesizes a node entry and preserves this counter rather than reporting zero. |
| `Metrics.retries` | Redeliveries. Always zero today: a failed object is not retried by an event, it waits for the scanner pass described below. Read `errors` instead. |
Healthy means: the same topology is visible on all sites, no pending operation,
no peer error, no failed retry escalation, required bucket/IAM state is in sync,
and queue/error counters are stable or falling. Counters are cumulative; alert on
their rate and on a backlog that does not drain, not merely on a non-zero total.
## Pending operations and recovery
`pendingOperation` contains `operation`, an opaque `id`, `pendingPeers`, and
`ackedPeers`. Do not edit the site-replication state object by hand. The marker
is the crash-recovery journal and removing it can make a partially applied
operation look complete.
The heavyweight reconciler runs once at startup and every 600 seconds. The
lightweight retry drain runs every 30 seconds. A restart is therefore a valid
way to cause an immediate heavyweight pass after the underlying fault has been
fixed, but it is not a substitute for fixing connectivity, credentials, TLS,
or the remote endpoint.
### `remove`
The original topology and each peer acknowledgement are persisted before the
operation finalizes. While peers remain in `pendingPeers`, restore access to
them and wait for reconciliation. If a peer is permanently gone, a new remove
request may remove all currently active unacknowledged peers; RustFS permits
that request and then finalizes against the remaining topology. Removing the
local site or all sites is also an explicit completion path.
Do not re-add a site merely to hide this marker. First compare the topology on
all reachable peers. If the same operation ID makes no progress for more than
one heavyweight interval, collect `PeerErrors`, `RetryStats`, and the
site-replication logs before retrying the remove.
### `rotate-svc-acct`
Service-account rotation keeps the candidate secrets and peer acknowledgements
until every current remote peer accepts the rotation. Restore the failing peer
and allow the reconciler to resume it. Do not manually delete either candidate
credential during this window: doing so can remove the only credential that a
not-yet-acknowledged peer accepts.
After the marker clears, verify `replicate status` from every site, then retire
any separately retained old credential material according to local policy.
### `endpoint-refresh`
An endpoint, CA, or TLS-verification edit first refreshes the replication
target on every active peer and records acknowledgements. On startup and every
heavyweight pass, RustFS probes peer capability, uses the endpoint-refresh API
when supported (or the legacy peer-edit fallback), refreshes local bucket
targets, and commits the edit only after every still-active peer acknowledges.
If this marker is stuck:
1. Confirm that the proposed endpoint and CA are correct and reachable from
every site, not only from the admin workstation.
2. Restore the site-replication service account and TLS trust path.
3. Wait for one 600-second pass or restart one healthy node to trigger the
startup pass.
4. Re-run the identical edit only if the operation remains visible; a different
endpoint edit is rejected while the existing refresh is pending. The journal
pins the edit's payload, so a re-run without `--replicate-ilm-expiry` keeps
the value the first attempt recorded, and a re-run asking for a different
value is rejected. Finish or remove the pending refresh before changing it.
A peer removed from the topology no longer blocks completion. A remove request
is accepted when it removes every active unacknowledged peer.
While this marker is present, control-plane retry replay to the other peers
keeps running, but bucket wiring reconciliation waits: it rewrites the same
targets the refresh is changing. Expect bucket-level drift on this site to
persist until the refresh settles.
## Outage recovery and convergence time
Control-plane retry begins on the 30-second drain, while heavyweight snapshots,
pending topology operations, and bucket wiring are revisited on the 600-second
pass. Object MRF entries are persisted every 10 seconds by default and target
health is probed every 5 seconds. These are scheduling bounds, not delivery
SLAs: network timeouts and the amount of queued work add to them.
Objects that must be rediscovered by the scanner have this conservative upper
bound before discovery:
```text
RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES
× max(RUSTFS_SCANNER_CYCLE, actual duration of one scanner cycle)
```
The defaults re-descend a compacted directory every 16 cycles. A practical
production starting point for a tighter recovery objective is
`RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES=4`; `1` forces re-descent every cycle.
Measure the additional disk and metadata load before lowering it further or
tuning the scanner cadence. For an immediate operator-driven recovery, start a
site resync with `mc admin replicate resync start` and monitor its status.
Transfer time after discovery remains proportional to backlog size, bandwidth,
worker capacity, and target latency. Use queue depth and the rate of
`Metrics.errors` rather than the formula alone to decide whether convergence is
progressing.
## Pairing sites that already contain data
When more than one requested site is non-empty, preflight considers each bucket
name held by more than one site:
- versioning must be `Enabled` on every site holding the shared bucket;
- Object Lock enablement must be identical on every holder.
A bucket present on only one site is safe: post-add backfill creates it on the
other peers. A shared unversioned bucket is rejected because merging can
overwrite the only copy of an object. An Object Lock mismatch is rejected
because lock enablement cannot be changed after bucket creation and convergence
could otherwise strip a WORM guarantee.
If preflight rejects the pair, keep the authoritative copy, delete the
conflicting bucket (or its contents) from all other sites, run `replicate add`
again, and then start `replicate resync` from the surviving site. Back up and
validate the authoritative data before deleting anything.
## IAM convergence and repair boundary
Ordinary IAM changes are delivered to each peer. A successful bulk IAM import
also schedules one collapsed full-IAM snapshot per remote peer. A failed IAM
deletion is replayed before that snapshot so the snapshot cannot re-create a
principal or grant that was already revoked.
The safety state has two bounds:
- deletion high-water marks are retained for 30 days;
- deletion replay bodies are capped at 256 distinct entities per peer.
Repeated deletion of the same entity replaces its saved body. When the per-peer
cap is exceeded or the body cannot be serialized, the retry entry remains
escalated rather than pretending the deletion is replayable. An item from an
older sender without a source timestamp cannot install the 30-day high-water
mark, so verify it explicitly after a prolonged split. A successful drain
clears replay bodies; removing the peer prunes its bodies. For an escalated IAM
retry, use the site-replication repair workflow for the affected peer and IAM
family, then verify users, service accounts, groups, policies, and mappings on
both sides. Repair is the operator's explicit accountability transfer and
clears the saved deletion bodies only after the IAM repair succeeds.
A group's status converges in one direction. An explicit disable is applied
everywhere, including through a snapshot, but a membership change never
carries an enable - it would otherwise re-enable a group frozen on the
receiving site. If a group ended up disabled on one site only, re-enable it
there explicitly with `mc admin group enable`; a snapshot or repair will not
do it.
Treat IAM divergence as a security incident: a user deleted on one site can
remain usable on an unreachable peer until replay or repair completes. A peer
whose IAM entry is escalated does not receive scheduled snapshots either -
including the one a bulk import schedules - until the repair settles it.
## Encrypted objects
| Source form | Replication behavior | Fail-closed condition |
|---|---|---|
| SSE-S3 | The source decrypts the object; the request sends only `AES256` intent; the destination encrypts with its own KMS. Source envelope material never leaves the site. | The destination cannot satisfy the encryption request, or the source metadata is incomplete/unsupported. The replica is `FAILED`; plaintext is not silently stored. |
| SSE-KMS | The source decrypts the object; the request sends `aws:kms` intent without the source-local key ID; the destination selects its own configured KMS key. | Either side cannot decrypt/encrypt, or the metadata mixes incompatible encryption evidence. |
| SSE-C | Stored ciphertext and the required SSE-C replication transport metadata pass through. RustFS verifies target evidence before accepting the replica. | The target does not echo the customer-algorithm evidence, required material/layout is absent, or the metadata is ambiguous. |
Unknown MinIO/RustFS encryption markers are never forwarded as ordinary user
metadata. They fail replication so an operator must migrate or repair the
object with a supported format.
## Rolling upgrades and rollback
Keep every node in one site on the same version whenever possible. Upgrade all
nodes of one site consecutively, verify its startup reconciliation and status,
then move to the next site. Do not intentionally leave a site mixed-version:
admin requests can land on different nodes, and an older node may not resume a
new pending-operation shape or expose its health fields.
Current state additions are optional and defaulted, so older readers ignore
them. The target-version ledger is stored as dual-prefixed internal object
metadata and is also ignored by older readers; rollback does not corrupt the
object format, but older code loses the assigned-version routing improvement.
Before rolling back across the fix that retains the data directory of a version
awaiting purge replication (rustfs/rustfs#7307), ensure no version purge is
pending. Older code can free that retained version's data directory before the
remote purge is acknowledged, leaving unreadable metadata and blocking bucket
deletion. Drain or repair replication and take a metadata/data backup first.
## Runtime knobs
These values are read when the owning background task starts. Restart the
server after changing them. The millisecond intervals have a 10 ms floor;
invalid values fall back to the default with a warning.
| Variable | Default | Effect |
|---|---:|---|
| `RUSTFS_REPL_HEALTH_CHECK_INTERVAL_MS` | `5000` | Remote-target health probe interval. Lowering it increases outbound probes. |
| `RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS` | `10000` | Maximum periodic interval between MRF persistence flushes; 1,000 new entries also trigger a flush. |
| `RUSTFS_REPL_RESYNC_POLL_MAX_MS` | `60000` | Upper bound for randomized resync retry-poll sleep. |
| `RUSTFS_REPL_RESYNC_MAX_JOBS` | `2` | Concurrent resync jobs; values are bounded to `1..=32`. |
Transport-specific controls and target behavior are documented in
[Replication outbound transport](replication-outbound-transport.md). Validate a
new destination with [Replication target check](replication-check.md), and read
[Replication object size limits](replication-object-size-limits.md) before
moving large objects.
+52
View File
@@ -627,6 +627,32 @@ pub(crate) async fn cluster_replication_stats(bucket: &str, context: Option<Arc<
.await
}
/// Reload the bucket's metadata on every peer so a follow-up
/// `put-bucket-replication` on another node does not read a stale target.
///
/// Best effort, like every S3 bucket-config write path
/// (`app::bucket_usecase::notify_bucket_metadata_reload`): the target is
/// already persisted and live on this node, and the 15-minute refresh closes
/// the gap, so a peer that cannot be reached must not turn a completed write
/// into a failed request.
async fn notify_remote_target_metadata_reload(bucket: &str, context: Option<Arc<AppContext>>, action: &'static str) {
let Some(notification_system) = current_notification_system_for_context(context.as_deref()) else {
return;
};
if let Err(err) = notification_system.load_bucket_metadata(bucket).await {
warn!(
event = EVENT_ADMIN_REMOTE_TARGET_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_REPLICATION,
action = action,
result = "peer_metadata_reload_failed",
bucket = %bucket,
error = ?err,
"admin remote target state"
);
}
}
fn unique_replication_peers(peer_clients: &[Option<PeerRestClient>]) -> (Vec<&PeerRestClient>, u32) {
let mut seen_grid_hosts = HashSet::new();
let peers: Vec<_> = peer_clients
@@ -699,6 +725,7 @@ pub struct SetRemoteTargetHandler {}
impl Operation for SetRemoteTargetHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let cred = validate_replication_admin_request(&req, AdminAction::SetBucketTargetAction).await?;
let app_context = app_context_from_req(&req);
let queries = extract_query_params(&req.uri);
@@ -926,6 +953,8 @@ impl Operation for SetRemoteTargetHandler {
.map_err(map_bucket_target_error)?;
let _targets_guard = lock_bucket_targets_metadata(bucket).await;
let arn = persist_remote_target_write(bucket, remote_target, incarnation, mode).await?;
drop(_targets_guard);
notify_remote_target_metadata_reload(bucket, app_context, "set_remote_target").await;
let arn_str = serde_json::to_string(&arn)
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize target ARN"))?;
@@ -1006,6 +1035,7 @@ pub struct RemoveRemoteTargetHandler {}
impl Operation for RemoveRemoteTargetHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
validate_replication_admin_request(&req, AdminAction::SetBucketTargetAction).await?;
let app_context = app_context_from_req(&req);
debug!("remove remote target called");
let queries = extract_query_params(&req.uri);
@@ -1081,6 +1111,7 @@ impl Operation for RemoveRemoteTargetHandler {
}
let json_targets = serde_json::to_vec(&targets)
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize targets"))?;
let notification_bucket = bucket.clone();
let bucket = bucket.clone();
let arn = arn_str.clone();
// The pool cancellation owns a detached task. Both outer guards must
@@ -1101,6 +1132,8 @@ impl Operation for RemoveRemoteTargetHandler {
S3Error::with_message(S3ErrorCode::InternalError, format!("remote target removal task failed: {error}"))
})??;
notify_remote_target_metadata_reload(&notification_bucket, app_context, "remove_remote_target").await;
Ok(S3Response::new((StatusCode::NO_CONTENT, Body::from("".to_string()))))
}
}
@@ -1787,6 +1820,25 @@ mod tests {
pairs.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect()
}
#[test]
fn remote_target_writes_notify_peer_metadata_caches() {
let source = include_str!("replication.rs");
for (start, end) in [
("impl Operation for SetRemoteTargetHandler", "pub struct ListRemoteTargetHandler"),
("impl Operation for RemoveRemoteTargetHandler", "async fn cancel_active_resync_intent"),
] {
let body = source
.split(start)
.nth(1)
.and_then(|rest| rest.split(end).next())
.expect(start);
assert!(
body.contains("notify_remote_target_metadata_reload"),
"{start} must notify every node before returning success"
);
}
}
#[test]
fn update_ops_parse_minio_query_contract() {
let ops = parse_remote_target_update_ops(&query_map(&[
File diff suppressed because it is too large Load Diff
+28
View File
@@ -1318,6 +1318,24 @@ impl Operation for ImportIam {
failed,
};
// The entities are already imported locally. A snapshot that cannot be
// scheduled is a convergence delay the reconcile pass still closes, so
// it must not turn a completed import into a failed request - the same
// best-effort contract every other site-replication hook here follows.
if let Err(err) =
crate::site_replication::enqueue_site_replication_iam_snapshot("iam import scheduled a full snapshot").await
{
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_USER,
event = EVENT_ADMIN_USER_STATE,
action = "import_iam",
result = "site_replication_snapshot_not_scheduled",
error = ?err,
"admin user state"
);
}
let body = serde_json::to_vec(&ret).map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, e.to_string()))?;
let mut header = HeaderMap::new();
@@ -1424,6 +1442,16 @@ mod tests {
assert!(include_str!("user.rs").contains(mapper_call));
}
#[test]
fn import_iam_enqueues_a_site_replication_snapshot() {
let body = source_block(include_str!("user.rs"), "impl Operation for ImportIam");
assert!(
body.contains("enqueue_site_replication_iam_snapshot"),
"a successful IAM import must schedule a full IAM snapshot for every remote site"
);
}
#[test]
fn test_should_check_deny_only_for_regular_self_request() {
let cred = Credentials {
@@ -409,6 +409,26 @@ fn transfer_summaries(stats: &InternalReplicationStats) -> (XferSummaryWire, Tar
(summary, per_target)
}
/// Node-level failure counters for `errors`. The sibling `retries` field
/// stays zero on purpose: it means redeliveries in the minio-go shape, and a
/// failed object is not retried by an event today (it waits for the scanner's
/// heal pass), so reporting failures there would claim a redelivery that
/// never happened.
fn failure_counters(stats: &InternalReplicationStats) -> CounterSummaryWire {
let (total, last1m, last1hr) = stats.stats.values().fold((0i64, 0i64, 0i64), |acc, stat| {
(
acc.0.saturating_add(stat.fail_stats.count),
acc.1.saturating_add(stat.fail_stats.last_minute.count),
acc.2.saturating_add(stat.fail_stats.last_hour.count),
)
});
CounterSummaryWire {
total: u64::try_from(total.max(0)).unwrap_or_default(),
last1m: u64::try_from(last1m.max(0)).unwrap_or_default(),
last1hr: u64::try_from(last1hr.max(0)).unwrap_or_default(),
}
}
impl MetricsV2Wire {
/// Project the aggregated internal stats onto the `MetricsV2` shape.
///
@@ -418,6 +438,7 @@ impl MetricsV2Wire {
/// `queueStats.nodes` and treats an empty list as "no data".
pub(crate) fn from_stats(bucket_stats: &BucketStats, node_name: &str) -> Self {
let (xfer_stats, tgt_xfer_stats) = transfer_summaries(&bucket_stats.replication_stats);
let failed = failure_counters(&bucket_stats.replication_stats);
let mut nodes: Vec<ReplQNodeStatsWire> = bucket_stats
.queue_stats
.nodes
@@ -436,6 +457,7 @@ impl MetricsV2Wire {
q_stats: InQueueMetricWire::from(&bucket_stats.replication_stats.q_stat),
xfer_stats: xfer_stats.clone(),
tgt_xfer_stats: tgt_xfer_stats.clone(),
errors: failed,
..Default::default()
});
} else {
@@ -444,6 +466,7 @@ impl MetricsV2Wire {
if let Some(first) = nodes.first_mut() {
first.xfer_stats = xfer_stats.clone();
first.tgt_xfer_stats = tgt_xfer_stats.clone();
first.errors = failed;
}
}
@@ -478,6 +501,12 @@ mod tests {
target.replicated_size = 4096;
target.failed.count = 3;
target.failed.size = 900;
target.fail_stats.count = 3;
target.fail_stats.size = 900;
target.fail_stats.last_minute.count = 2;
target.fail_stats.last_minute.size = 600;
target.fail_stats.last_hour.count = 3;
target.fail_stats.last_hour.size = 900;
target.bandwidth_limit_bytes_per_sec = 1024;
target.current_bandwidth_bytes_per_sec = 512.5;
stats
@@ -537,6 +566,10 @@ mod tests {
assert_eq!(node["queueStats"]["peak"], node["queueStats"]["max"]);
assert!(node["activeWorkers"].get("curr").is_some());
assert!(node["transferSummary"].get("Total").is_some());
assert_eq!(node["errors"]["total"], 3);
assert_eq!(node["errors"]["last1m"], 2);
assert_eq!(node["errors"]["last1hr"], 3);
assert_eq!(node["retries"]["total"], 0, "failures are not redeliveries; retries must not claim one");
assert_eq!(json["downtimeInfo"], serde_json::json!({}));
}
+86 -1
View File
@@ -217,6 +217,26 @@ pub(crate) fn settle_observed_site_replication_retry_event(
before.saturating_sub(queue.len())
}
/// Make sure `peer` has a collapsed entry for `path` without counting the
/// call as a delivery failure. A bulk local mutation (`import-iam`) needs the
/// entry to exist so the next drain sends the snapshot; routing it through
/// [`upsert_site_replication_retry_event`] would raise `retry_count` on every
/// import and escalate a healthy peer to `failed` after
/// [`SITE_REPLICATION_RETRY_FAILED_AFTER`] of them, with the scheduling note
/// shown to operators as `lastError`.
pub(crate) fn ensure_site_replication_retry_event(
queue: &mut Vec<SiteReplicationRetryEvent>,
peer: &PeerInfo,
path: &str,
reason: &str,
) -> S3Result<Vec<SiteReplicationRetryEvent>> {
let path = collapsed_retry_queue_path(path).unwrap_or(path);
if queue.iter().any(|event| retry_event_matches(event, peer, path)) {
return Ok(Vec::new());
}
push_site_replication_retry_event(queue, peer, path, summarize_peer_error_detail(reason), false, None)
}
pub(crate) fn upsert_site_replication_retry_event(
queue: &mut Vec<SiteReplicationRetryEvent>,
peer: &PeerInfo,
@@ -244,6 +264,17 @@ pub(crate) fn upsert_site_replication_retry_event(
return Ok(Vec::new());
}
push_site_replication_retry_event(queue, peer, path, detail, peer_unreachable, generation)
}
fn push_site_replication_retry_event(
queue: &mut Vec<SiteReplicationRetryEvent>,
peer: &PeerInfo,
path: &str,
detail: String,
peer_unreachable: bool,
generation: Option<u64>,
) -> S3Result<Vec<SiteReplicationRetryEvent>> {
let slots_needed = queue
.len()
.saturating_add(1)
@@ -274,7 +305,7 @@ pub(crate) fn upsert_site_replication_retry_event(
retry_count: 1,
failed: false,
last_error: detail,
updated_at: Some(now),
updated_at: Some(OffsetDateTime::now_utc()),
edit_generation: generation,
peer_unreachable,
deletions_recorded: false,
@@ -365,6 +396,60 @@ pub(crate) async fn enqueue_site_replication_retry_event_for_generation(
}
}
/// Returns the number of peers whose snapshot entry is escalated and therefore
/// will not carry this scheduling: the marker records a deletion that a
/// snapshot cannot replay, and only a repair settles it, so clearing it to make
/// the entry drainable again would drop that liability.
pub(crate) fn record_iam_snapshot_retries(
state: &mut SiteReplicationState,
local_peer: &PeerInfo,
reason: &str,
) -> S3Result<usize> {
let peers = state
.peers
.values()
.filter(|peer| {
peer.deployment_id != local_peer.deployment_id && !same_identity_endpoint(&peer.endpoint, &local_peer.endpoint)
})
.cloned()
.collect::<Vec<_>>();
let mut escalated = 0usize;
for peer in peers {
if state.retry_queue.iter().any(|event| {
retry_event_matches(event, &peer, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH)
&& event.last_error == SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER
}) {
escalated += 1;
continue;
}
ensure_site_replication_retry_event(&mut state.retry_queue, &peer, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH, reason)?;
}
Ok(escalated)
}
/// Schedule one collapsed full-IAM snapshot per remote peer after a bulk
/// local mutation such as `import-iam`.
pub(crate) async fn enqueue_site_replication_iam_snapshot(reason: &str) -> S3Result<()> {
let state = load_site_replication_state().await?;
if !state.enabled() {
return Ok(());
}
let local_peer = current_local_runtime_peer(&state);
let reason = reason.to_string();
let escalated = update_site_replication_state(move |state| record_iam_snapshot_retries(state, &local_peer, &reason)).await?;
if escalated > 0 {
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
escalated,
result = "iam_snapshot_not_scheduled_for_escalated_peer",
"site replication peers hold an escalated IAM entry; the snapshot waits for a repair"
);
}
Ok(())
}
pub(crate) const SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH: &str = "/rustfs/admin/v3/site-replication/peer/iam-item";
/// Per-peer cap on recorded deletion bodies. Beyond it the peer's collapsed
+2
View File
@@ -168,6 +168,8 @@ mod rfc3339_map {
pub(crate) struct PendingEndpointRefresh {
pub(crate) id: String,
pub(crate) peer: PeerInfo,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) ilm_expiry_override: Option<bool>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub(crate) remote_peers: BTreeMap<String, PeerInfo>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
+116
View File
@@ -693,6 +693,122 @@ fn test_record_iam_deletion_marks_newest_wins_and_expires_by_age_only() {
);
}
/// Scheduling a snapshot is not a delivery failure. Repeated imports - the
/// normal way a bulk IAM migration is done, one archive at a time - must not
/// walk the peer's entry up to the escalation threshold and report a healthy
/// site as `retryStats.failed` with the scheduling note as its `lastError`.
#[test]
fn repeated_iam_import_snapshots_do_not_escalate_a_healthy_peer() {
let local = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let remote = PeerInfo {
deployment_id: "remote-a".to_string(),
..peer("remote-a", "https://a.example.com")
};
let mut state = SiteReplicationState {
peers: BTreeMap::from([
(local.deployment_id.clone(), local.clone()),
(remote.deployment_id.clone(), remote),
]),
..Default::default()
};
for _ in 0..(SITE_REPLICATION_RETRY_FAILED_AFTER + 2) {
record_iam_snapshot_retries(&mut state, &local, "iam import scheduled a full snapshot").expect("record snapshot");
}
assert_eq!(state.retry_queue.len(), 1);
let event = &state.retry_queue[0];
assert_eq!(event.retry_count, 1, "a schedule must not count as a delivery attempt");
assert!(!event.failed, "a scheduled snapshot must not report as an escalated failure");
}
/// An escalated entry records a deletion a snapshot cannot replay: only a
/// repair settles it. Scheduling an import snapshot must not clear that
/// marker to make the entry drainable again, and the peer it skips has to be
/// reported rather than silently left behind.
#[test]
fn an_escalated_peer_keeps_its_marker_and_is_reported() {
let local = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let remote = PeerInfo {
deployment_id: "remote-a".to_string(),
..peer("remote-a", "https://a.example.com")
};
let mut state = SiteReplicationState {
peers: BTreeMap::from([
(local.deployment_id.clone(), local.clone()),
(remote.deployment_id.clone(), remote.clone()),
]),
retry_queue: vec![SiteReplicationRetryEvent {
id: "escalated".to_string(),
peer_deployment_id: remote.deployment_id.clone(),
peer_endpoint: remote.endpoint,
path: SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH.to_string(),
retry_count: SITE_REPLICATION_RETRY_FAILED_AFTER,
failed: true,
last_error: SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER.to_string(),
deletions_recorded: true,
..Default::default()
}],
..Default::default()
};
let escalated =
record_iam_snapshot_retries(&mut state, &local, "iam import scheduled a full snapshot").expect("record snapshot retries");
assert_eq!(escalated, 1);
assert_eq!(state.retry_queue.len(), 1);
assert_eq!(
state.retry_queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER,
"the unreplayable-deletion marker must survive a snapshot schedule"
);
}
#[test]
fn iam_import_snapshot_retry_is_recorded_once_per_remote_peer() {
let local = PeerInfo {
deployment_id: "local-dep".to_string(),
..peer("local", "https://local.example.com")
};
let remote_a = PeerInfo {
deployment_id: "remote-a".to_string(),
..peer("remote-a", "https://a.example.com")
};
let remote_b = PeerInfo {
deployment_id: "remote-b".to_string(),
..peer("remote-b", "https://b.example.com")
};
let mut state = SiteReplicationState {
peers: BTreeMap::from([
(local.deployment_id.clone(), local.clone()),
(remote_a.deployment_id.clone(), remote_a),
(remote_b.deployment_id.clone(), remote_b),
]),
..Default::default()
};
record_iam_snapshot_retries(&mut state, &local, "IAM import snapshot pending").expect("record snapshot retries");
assert_eq!(state.retry_queue.len(), 2);
assert!(
state
.retry_queue
.iter()
.all(|event| event.path == SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH)
);
assert!(
state
.retry_queue
.iter()
.all(|event| event.peer_deployment_id != local.deployment_id)
);
}
/// A failed deletion delivery persists a replay record next to the collapsed
/// retry entry; a fresh entry is stamped `deletions_recorded` so a later
/// replay can settle it, and a repeated deletion of the same entity keeps the
+69 -4
View File
@@ -17,6 +17,7 @@
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::sync::{Arc, LazyLock};
use std::time::Duration;
use rand::RngExt as _;
use rustfs_storage_api as storage_contracts;
@@ -836,6 +837,39 @@ impl StorageReplicationStatsHandle {
pub(crate) async fn site_metrics_snapshot(&self) -> ReplicationSiteMetricsSnapshot {
let metrics = self.inner.get_sr_metrics_for_node().await;
// Aggregate under the read lock rather than through `get_all`: that
// clones every bucket's stats, and `FailStats.recent` is bounded only
// by the one-hour window, so an unreachable target under load - the
// very case an operator polls this for - makes the copy large. The
// windows come from the live samples; the serialized `last_minute` /
// `last_hour` snapshots are stamped onto per-bucket clones elsewhere
// and stay zero in this node-local cache.
let (
failed_count,
failed_bytes,
failed_last_minute_count,
failed_last_minute_bytes,
failed_last_hour_count,
failed_last_hour_bytes,
) = {
let cache = self.inner.cache.read().await;
cache
.values()
.flat_map(|bucket| bucket.stats.values())
.fold((0i64, 0i64, 0i64, 0i64, 0i64, 0i64), |totals, stat| {
let (minute, hour) = stat
.fail_stats
.recent_windows(Duration::from_secs(60), Duration::from_secs(3600));
(
totals.0.saturating_add(stat.fail_stats.count),
totals.1.saturating_add(stat.fail_stats.size),
totals.2.saturating_add(minute.count),
totals.3.saturating_add(minute.size),
totals.4.saturating_add(hour.count),
totals.5.saturating_add(hour.size),
)
})
};
ReplicationSiteMetricsSnapshot {
uptime: metrics.uptime,
queued_curr_count: metrics.queued.curr.count,
@@ -859,6 +893,12 @@ impl StorageReplicationStatsHandle {
proxy_delete_tag_failed: metrics.proxied.delete_tag_failed,
replica_size: metrics.replica_size,
replica_count: metrics.replica_count,
failed_count,
failed_bytes,
failed_last_minute_count,
failed_last_minute_bytes,
failed_last_hour_count,
failed_last_hour_bytes,
}
}
@@ -899,6 +939,12 @@ pub(crate) struct ReplicationSiteMetricsSnapshot {
pub(crate) proxy_delete_tag_failed: i64,
pub(crate) replica_size: i64,
pub(crate) replica_count: i64,
pub(crate) failed_count: i64,
pub(crate) failed_bytes: i64,
pub(crate) failed_last_minute_count: i64,
pub(crate) failed_last_minute_bytes: i64,
pub(crate) failed_last_hour_count: i64,
pub(crate) failed_last_hour_bytes: i64,
}
pub(crate) async fn get_local_server_property() -> rustfs_madmin::ServerProperties {
@@ -2043,13 +2089,32 @@ pub(crate) async fn init_compression_total_memory_from_backend(store: Arc<ECStor
#[cfg(test)]
mod tests {
use super::{
BUCKET_RESYNC_LOCK_RETRY_MAX_MS, apply_active_resync_intents, bucket_resync_transaction_lock_retry_ceiling_ms,
bucket_resync_transaction_lock_retry_delay, bucket_resync_transaction_lock_retry_reason,
bucket_targets_metadata_lock_shard, ecstore_bucket, lock_bucket_targets_metadata, new_instance_ctx,
retry_bucket_resync_transaction_lock, scanner_maintenance_config_file,
BUCKET_RESYNC_LOCK_RETRY_MAX_MS, StorageReplicationStatsHandle, apply_active_resync_intents,
bucket_resync_transaction_lock_retry_ceiling_ms, bucket_resync_transaction_lock_retry_delay,
bucket_resync_transaction_lock_retry_reason, bucket_targets_metadata_lock_shard, ecstore_bucket,
lock_bucket_targets_metadata, new_instance_ctx, retry_bucket_resync_transaction_lock, scanner_maintenance_config_file,
};
use std::time::Duration;
#[tokio::test]
async fn site_metrics_snapshot_includes_live_failure_windows() {
let stats = StorageReplicationStatsHandle::new();
let mut target = ecstore_bucket::replication::BucketReplicationStat::default();
target.fail_stats.add_size(2048, None::<&std::io::Error>);
let mut bucket = ecstore_bucket::replication::BucketReplicationStats::new();
bucket.stats.insert("arn:replication::remote:photos".to_string(), target);
stats.inner.cache.write().await.insert("photos".to_string(), bucket);
let snapshot = stats.site_metrics_snapshot().await;
assert_eq!(snapshot.failed_count, 1);
assert_eq!(snapshot.failed_bytes, 2048);
assert_eq!(snapshot.failed_last_minute_count, 1);
assert_eq!(snapshot.failed_last_minute_bytes, 2048);
assert_eq!(snapshot.failed_last_hour_count, 1);
assert_eq!(snapshot.failed_last_hour_bytes, 2048);
}
#[tokio::test]
async fn bucket_target_metadata_locks_serialize_only_matching_shards() {
let bucket = "bucket-target-lock";
+4 -1
View File
@@ -55,8 +55,11 @@ cd "$(dirname "$0")/.."
# now reports an unreadable configuration as a plain string instead of raising
# an S3 error per arm (24 invocation lines removed from
# rustfs/src/admin/handlers/bucket_meta.rs; measured after merging the two).
# 1589 -> 1588 on 2026-09-08: the GA blocker set (rustfs/backlog#2366) added
# three invocation lines to the endpoint-refresh paths and folded the five
# copies of the concurrent-change error into one constructor, netting -1.
S3S_IMPORT_FILES_BASELINE=213
S3_ERROR_LINES_BASELINE=1589
S3_ERROR_LINES_BASELINE=1588
# ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not
# know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*
# client was extracted to crates/s3-client, where s3s usage is legitimate;
+49
View File
@@ -837,6 +837,55 @@ emit_step_result() {
self.assertIn(value, contents)
self.assertNotIn("OLD RUN EVIDENCE", contents)
def test_performance_commands_bind_runner_selection_and_preserve_failures(self) -> None:
self.prepare("performance")
source = self.source.splitlines()
job = yaml_block(source, "performance-test", 2)
runner = WorkflowSteps()
runner.directory = self.directory / "workspace with spaces"
scripts = runner.directory / "auto-testing"
scripts.mkdir(parents=True)
wrapper = scripts / "rustfs_performance_test.sh"
wrapper.write_text(f"#!{sys.executable}\nimport json, os, sys\n" +
"print(json.dumps({'args': sys.argv[1:], 'env': {key: os.environ.get(key) for key in " +
"('RUSTFS_BENCH_SCRIPT', 'RUSTFS_WARP_METHODS', 'RUSTFS_WARP_SIZES', " +
"'RUSTFS_WARP_DURATION', 'RUSTFS_WARP_CONCURRENCY', 'WARP_METHODS', " +
"'WARP_SIZES', 'WARP_DURATION', 'WARP_CONCURRENCY')}}))\n" +
"sys.exit(int(os.environ['FAKE_BENCH_EXIT']))\n")
wrapper.chmod(0o755)
runner.steps = named_steps(job)
for methods, sizes, duration, concurrency in (
("get", "1KiB", "1s", "7"), ("all", "all", "5m", "64"), ("", "", "5m", "64")
):
runner.context = {"github.workspace": str(runner.directory), "inputs.test_method": methods,
"inputs.object_size": sizes, "inputs.warp_duration || '5m'": duration,
"inputs.warp_concurrency || '64'": concurrency}
runner.env = {**self.env, "RUSTFS_BENCH_SCRIPT": "/unverified/home-script.sh",
"RUSTFS_WARP_METHODS": "put", "RUSTFS_WARP_SIZES": "64MiB",
"RUSTFS_WARP_DURATION": "99h", "RUSTFS_WARP_CONCURRENCY": "2",
"WARP_DURATION": "88h", "WARP_CONCURRENCY": "3", "WARP_METHODS": "mixed", "WARP_SIZES": "32MiB",
"LOG_FILE": str(self.directory / "suite.log")}
runner.env.update(runner.step_env(job, indent=4))
for step, number in (("Run benchmark (GET/PUT/MIXED)", "5"), ("Analyze results", "6")):
for code in (0, 42):
with self.subTest(methods=methods, sizes=sizes, step=step, exit=code):
runner.env["FAKE_BENCH_EXIT"] = str(code)
result = runner.run_step(step)
self.assertEqual(result.returncode, code, result.stderr)
invocation = json.loads(result.stdout)
expected = ["--step", number, "-y", "--log-file", runner.env["LOG_FILE"]]
self.assertEqual(invocation["args"], expected)
self.assertEqual(invocation["env"]["RUSTFS_BENCH_SCRIPT"], str(scripts / "rustfs_performance_testing.sh"))
self.assertEqual(invocation["env"]["RUSTFS_WARP_METHODS"], methods)
self.assertEqual(invocation["env"]["RUSTFS_WARP_SIZES"], sizes)
self.assertEqual(invocation["env"]["RUSTFS_WARP_DURATION"], duration)
self.assertEqual(invocation["env"]["RUSTFS_WARP_CONCURRENCY"], concurrency)
if number == "6":
self.assertEqual(invocation["env"]["WARP_METHODS"], methods)
self.assertEqual(invocation["env"]["WARP_SIZES"], sizes)
self.assertEqual(invocation["env"]["WARP_DURATION"], duration)
self.assertEqual(invocation["env"]["WARP_CONCURRENCY"], concurrency)
if __name__ == "__main__":
unittest.main()