Compare commits

...

10 Commits

Author SHA1 Message Date
唐小鸭 6c37ebb951 fix(replication): address a replicated marker purge by the target's version id (backlog#2290)
A source-side DELETE ?versionId=<marker> replicates as a version purge, and
replicate_delete_to_target addressed it by the SOURCE marker id on every
target. A generic S3 target answers a DELETE of an unknown versionId with
204 and keeps its marker, so the purge reported success and the marker
stayed; the same event also spawned a second delayed-purge watcher that
journaled a duplicate intent. Real VMs (R6.1 in backlog#2080) failed on the
persisted-id fix alone because this path never consulted the mapping.

Resolve the target version through the recorded mapping for marker purges
(a corrupt record refuses, as the watcher does; nothing recorded keeps the
source-derived id for id-mirroring peers), and do not spawn the delayed
watcher for a version purge — that purge is the replication itself and its
failures reach the journal as a purge entry.
2026-09-05 19:40:40 +08:00
唐小鸭 850639c67f fix(site-replication): hash repair tasks and retry snapshots with sorted JSON keys (backlog#2289)
Service-account items carry their claims in a HashMap, and serde_json is
built with preserve_order, so two serializations of the same plan could
differ in key order. The repair preflight token then went stale between
dry-run and execute (412 on the real VMs once snapshots carried service
accounts) and a retry snapshot resend could never look stable. Serialize
through a key-sorted JSON value for the task id and the fingerprint.
2026-09-05 19:31:25 +08:00
唐小鸭 9032adcb12 fix(site-replication): judge service-account items against deletion marks too (backlog#2291)
The service-account receive path used only the live record's timestamp; a
deleted account left nothing to compare against, so a stale create from a
snapshot or a delayed delivery could recreate it. Consult the recorded
deletion mark when the record is absent, as the user path does. The site
replicator account is managed by join/rotate and stays exempt.
2026-09-05 18:45:59 +08:00
唐小鸭 18670b269b fix(replication): persist target delete-marker version ids in MRF purge intents (backlog#2290)
A delete-marker purge intent that outlived its watch window was journaled
without the version ids the targets assigned to the replicated markers.
Replay rebuilt the replication state from a blank ObjectInfo, so
`delete_marker_purge_version_id` fell back to the source marker id; a
generic S3 target that mints its own ids answers that DELETE with 204,
the entry was acknowledged and the real marker stayed on the target.

- `MrfReplicateEntry` gains `targetDeleteMarkerVersionIDs` (per-ARN map)
  and `targetDeleteMarkerVersionIDsCorrupt`; both default and are skipped
  when empty/false, so old journals decode to the pre-existing shape.
- `DeletedObjectReplicationInfo::to_mrf_entry` copies both from the
  source replication state; `reconstructed_heal_delete_info` restores
  them into the replayed state so the purge addresses the recorded id
  and a fail-closed refusal stays a refusal after restart.
- MRF envelope capability bit `TargetDeleteMarkerVersionIds` (1 << 4)
  fences the field like `DeleteMarkerMtime`; readers without the bit
  refuse envelopes that advertise it, current readers accept old ones.

(cherry picked from commit ddacaaa185fda7a5f426138ba5b179f986b862d9)
2026-09-05 18:45:59 +08:00
唐小鸭 c5e6b7259e fix(site-replication): stamp replicated bucket configs with the source updated_at (backlog#2292)
The bucket-meta receiver judged an incoming item stale by comparing its
source `updated_at` with the `*_config_updated_at` stamp of the config on
disk, but that stamp was the receiver's local clock at apply time
(`BucketMetadata::update_config`). A source edit newer than the applied one
but delivered after the local stamp was judged stale and acknowledged with
200: two quick source edits under delivery delay lose the second, and a peer
clock ahead of ours loses every follow-up edit inside the skew.

Add explicit-timestamp write entries, expanding rather than changing the
existing ones:

- `BucketMetadata::update_config_at`; `update_config` delegates to it with
  the local clock.
- `metadata_sys::update_if_incarnation_at`,
  `update_under_transaction_lock_at`, `update_quota_if_incarnation_at`,
  threaded through the shared write-guard path as `Option<OffsetDateTime>`
  (`None` keeps local stamping for every existing caller and for deletes).
- Re-exported through the ecstore `api` facade and the rustfs admin
  `storage_api::metadata_sys` facade.

`apply_bucket_meta_item` now persists policy, tags, versioning, object-lock,
sse, replication, quota and cors configs with the item's source time, so the
stored stamp equals the source `updatedAt` and staleness is judged source
time against source time. Items without `updated_at` keep the local stamp.
lc-config stays on the local stamp: its staleness axis is the in-document
`expiry_updated_at` the merge records, and the whole-config time only serves
as its deletion / legacy lower bound. Local (non-replicated) edits keep
stamping the local clock — they are the source.

(cherry picked from commit c1009c018b217ef9edc7773c8e56667ea7e77335)
2026-09-05 18:45:59 +08:00
唐小鸭 e834228926 fix(site-replication): keep IAM deletion marks so stale grants cannot follow a revoke (backlog#2291)
The staleness gate judged an incoming IAM item against the local record's
timestamp, but a full revoke deletes the record: `policy_db_set(.., "")`
removes the mapping, `delete_policy` the document, a user delete the
identity, and the IAM cache keeps only a per-entity watermark, no per-key
deletion time. With nothing left to compare against, a delayed older grant
was still applied after the revoke (real-VM case R6.3a of backlog#2080:
detach on A, revoke reaches B, an older mapping grant lands on B with 200
and re-grants access).

Keep a bounded, persisted map of deleted entity -> source `updatedAt` of the
newest deletion committed here in the site-replication state, written
through the state transaction in two places: the local IAM change hook
records the mark before broadcasting a deletion-shaped item, and the peer
item handler records it after applying (or idempotently no-op'ing) one. The
`policy`, `policy-mapping`, `group-info` and `iam-user` receive paths feed
that mark into the shared verdict when the record is absent, so a grant
older than the recorded deletion is acknowledged without being applied.
Group member removals are marked per member and a group delete marks the
group itself, so a stale re-add of a removed member is judged against the
newest of those marks.

Marks need a source timestamp: items without `updatedAt` (older peers) and
an unreadable state object fall back to today's behaviour and apply. The
map holds at most 1024 entries, evicting the oldest, and is cleared when
this site leaves the cluster.

(cherry picked from commit 0595c600d6091f583858176232d87ef2eed2bf8f)
2026-09-05 18:45:59 +08:00
唐小鸭 68499b6549 fix(site-replication): gate policy, mapping and group items on source updated_at (backlog#2291)
The `policy`, `policy-mapping` and `group-info` receive paths applied every
incoming item unconditionally, so a delayed older grant (wide policy body,
old mapping, old group add) overwrote a newer revoke on the peer. `iam-user`
and `service-account` already compared the item's `updatedAt` with the local
record.

Route the three paths through one pure verdict helper: an item older than
the local record is acknowledged without being applied; items without a
source timestamp and items targeting an absent record keep today's behaviour
(older peers, idempotent deletes from backlog#2071). Deletes are gated the
same way so an older delete cannot remove a newer record.

The group record's own timestamp now moves on every membership and status
change instead of staying at creation, so the gate judges group items
against the last change. Add the IamSys accessors the gate reads
(`get_policy_doc`, `get_mapped_policy_record`, `get_group_info`).

(cherry picked from commit 98c32093406cb47014b7eda2fe139f01079de337)
2026-09-05 18:45:59 +08:00
唐小鸭 2518fb5cd7 fix(site-replication): broadcast bucket ops to every peer and record each failure (backlog#2293)
The generic JSON broadcast (make/delete bucket, bucket-meta hook, bucket
ops) returned at the first failing peer, so peers later in deployment-id
order never received the request and got no retry event; a transport
construction failure recorded nothing at all.

Attempt every remote peer like the IAM change hook does: a success settles
the peer/path retry event, a failure (transport construction included)
enqueues one under the request path, and the first error is returned after
all peers were attempted.

(cherry picked from commit ce8f73bfd74bac61c383434780f27d17ca16d75e)
2026-09-05 18:45:59 +08:00
唐小鸭 b98b4821a8 fix(site-replication): run both make-bucket broadcast steps on every peer (backlog#2293)
With the generic broadcast now attempting every peer and returning the
first failure, stopping after the make step on that error skipped
configure-replication for the peers whose make had just succeeded, and no
retry event covered the gap. Run both steps and combine the results.

(cherry picked from commit cf5c0476dc4265b6214bf34c47a9ae4c61bf8e3e)
2026-09-05 18:45:59 +08:00
唐小鸭 b56ff5503f fix(site-replication): carry user credentials and service accounts in the IAM snapshot (backlog#2289)
The IAM snapshot used by the retry drain resend, repair and site-add
bootstrap was built from `list_users`, which strips secret keys and skips
service accounts. The plan builder dropped every user for lack of a
secret, so a user disable, secret rotation or service-account change
committed while a peer was unreachable never reached it — while the
collapsed retry entry was settled and repair reported success.

Read the credentials separately at plan time (`build_sr_iam_credentials`,
used only on peer-delivery paths) so `SRInfo`, which is served to admin
callers, stays secret-free. Users travel with secret, status and the user
record's own update time; service accounts (except the replicator's) travel
as the create item the live hook emits, after their parents. The receiver
applies a disabled status after creating a new service account, and the
retry snapshot tombstones removed service accounts like the other kinds.

`encode_service_account_replication_policy` moves into the infra layer so
the snapshot builder can share it with the live hook.

(cherry picked from commit bd8cd497c965ae331fafa20362763224cec3b5b2)
2026-09-05 18:45:59 +08:00
22 changed files with 2330 additions and 173 deletions
+3 -2
View File
@@ -204,8 +204,9 @@ pub mod bucket {
get_public_access_block_config, get_quota_config, get_replication_config, get_request_payment_config, get_sse_config,
get_tagging_config, get_versioning_config, get_website_config, init_bucket_metadata_sys, list_bucket_targets,
reload_bucket_metadata, remove_bucket_metadata, set_bucket_metadata, update,
update_bucket_targets_under_transaction_lock, update_config_with, update_if_incarnation, update_quota_if_incarnation,
update_under_transaction_lock,
update_bucket_targets_under_transaction_lock, update_config_with, update_if_incarnation, update_if_incarnation_at,
update_quota_if_incarnation, update_quota_if_incarnation_at, update_under_transaction_lock,
update_under_transaction_lock_at,
};
}
+47 -1
View File
@@ -802,9 +802,22 @@ impl BucketMetadata {
}
}
/// Replace one config payload and stamp its `*_config_updated_at` with the
/// local clock. This is the entry for edits that originate here: the
/// local write time is the edit's source time.
pub fn update_config(&mut self, config_file: &str, data: Vec<u8>) -> Result<OffsetDateTime> {
let updated = OffsetDateTime::now_utc();
self.update_config_at(config_file, data, OffsetDateTime::now_utc())
}
/// [`Self::update_config`] with an explicit `updated_at` stamp.
///
/// For a config replicated from another site the edit's source time is
/// the peer's `updated_at`, not the moment it lands here: staleness of
/// the next incoming item is judged against the stored stamp, so stamping
/// the local apply time would reject a newer source edit that was merely
/// delivered late (backlog#2292). Only replication receivers should pass
/// a foreign time; local edits keep [`Self::update_config`].
pub fn update_config_at(&mut self, config_file: &str, data: Vec<u8>, updated: OffsetDateTime) -> Result<OffsetDateTime> {
match config_file {
BUCKET_POLICY_CONFIG => {
self.policy_config_json = data;
@@ -1543,6 +1556,39 @@ mod test {
assert_eq!(metadata.bucket_incarnation_id, incarnation);
}
/// backlog#2292: a replicated config is stamped with the source
/// `updated_at` it was given, not the local clock, while the plain
/// `update_config` entry keeps stamping the local clock.
#[test]
fn update_config_at_stamps_the_given_time_and_update_config_stamps_now() {
let source_time = OffsetDateTime::now_utc() - time::Duration::hours(3);
let mut metadata = BucketMetadata::new("source-stamped");
let stamped = metadata
.update_config_at(BUCKET_POLICY_CONFIG, br#"{"Version":"2012-10-17","Statement":[]}"#.to_vec(), source_time)
.unwrap();
assert_eq!(stamped, source_time);
assert_eq!(metadata.policy_config_updated_at, source_time);
let tagging = b"<Tagging><TagSet><Tag><Key>k</Key><Value>v</Value></Tag></TagSet></Tagging>".to_vec();
let stamped = metadata
.update_config_at(BUCKET_TAGGING_CONFIG, tagging, source_time)
.unwrap();
assert_eq!(stamped, source_time);
assert_eq!(metadata.tagging_config_updated_at, source_time);
let before = OffsetDateTime::now_utc();
let stamped = metadata
.update_config(BUCKET_POLICY_CONFIG, br#"{"Version":"2012-10-17","Statement":[]}"#.to_vec())
.unwrap();
assert!(stamped >= before, "a local edit is stamped with the local clock");
assert_eq!(metadata.policy_config_updated_at, stamped);
assert_eq!(
metadata.tagging_config_updated_at, source_time,
"restamping one config must not move another config's stamp"
);
}
#[test]
fn object_locking_requires_lock_metadata_not_plain_versioning() {
use s3s::dto::ObjectLockEnabled;
+145 -14
View File
@@ -581,6 +581,32 @@ pub async fn update_if_incarnation(
config_file,
data,
Some(expected_incarnation_id),
None,
))
.await
}
/// [`update_if_incarnation`] stamping the config with `updated_at` instead of
/// the local clock.
///
/// For a site-replication receiver the edit's source time is the peer's
/// `updated_at`; persisting it keeps the stored `*_config_updated_at` on the
/// source clock so the next item's staleness is judged source-time against
/// source-time (backlog#2292). See [`BucketMetadata::update_config_at`].
pub async fn update_if_incarnation_at(
bucket: &str,
config_file: &str,
data: Vec<u8>,
expected_incarnation_id: Uuid,
updated_at: OffsetDateTime,
) -> Result<OffsetDateTime> {
Box::pin(update_with_sys_expected(
get_bucket_metadata_sys()?,
bucket,
config_file,
data,
Some(expected_incarnation_id),
Some(updated_at),
))
.await
}
@@ -612,18 +638,22 @@ async fn update_with_sys(
config_file: &str,
data: Vec<u8>,
) -> Result<OffsetDateTime> {
update_with_sys_expected(sys, bucket, config_file, data, None).await
update_with_sys_expected(sys, bucket, config_file, data, None, None).await
}
/// `updated_at` is the stamp persisted on the config; `None` uses the local
/// clock (the edit originates here), `Some` carries a replicated edit's
/// source time (backlog#2292).
async fn update_with_sys_expected(
sys: Arc<RwLock<BucketMetadataSys>>,
bucket: &str,
config_file: &str,
data: Vec<u8>,
expected_incarnation_id: Option<Uuid>,
updated_at: Option<OffsetDateTime>,
) -> Result<OffsetDateTime> {
let guard = acquire_config_write_guard_for_incarnation(sys.clone(), bucket, expected_incarnation_id).await?;
update_under_config_write_guard(sys, &guard, config_file, data).await
update_under_config_write_guard(sys, &guard, config_file, data, updated_at).await
}
/// [`delete`] against an explicitly supplied metadata system. See
@@ -786,7 +816,21 @@ pub async fn update_under_transaction_lock(
data: Vec<u8>,
) -> Result<OffsetDateTime> {
guard.ensure_valid(bucket)?;
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data).await
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data, None).await
}
/// [`update_under_transaction_lock`] stamping the config with `updated_at`
/// (a replicated edit's source time) instead of the local clock; see
/// [`update_if_incarnation_at`] (backlog#2292).
pub async fn update_under_transaction_lock_at(
guard: &BucketMetadataMutationGuard,
bucket: &str,
config_file: &str,
data: Vec<u8>,
updated_at: OffsetDateTime,
) -> Result<OffsetDateTime> {
guard.ensure_valid(bucket)?;
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data, Some(updated_at)).await
}
/// Clear one config file while the caller holds this bucket's transaction lock.
@@ -804,6 +848,29 @@ pub async fn update_quota_if_incarnation(
data: Vec<u8>,
expected_incarnation_id: Uuid,
proof: &crate::services::notification_sys::CrossPoolFenceFleetProofToken,
) -> Result<OffsetDateTime> {
update_quota_if_incarnation_stamped(bucket, data, expected_incarnation_id, proof, None).await
}
/// [`update_quota_if_incarnation`] stamping the quota config with
/// `updated_at` (a replicated edit's source time) instead of the local
/// clock; see [`update_if_incarnation_at`] (backlog#2292).
pub async fn update_quota_if_incarnation_at(
bucket: &str,
data: Vec<u8>,
expected_incarnation_id: Uuid,
proof: &crate::services::notification_sys::CrossPoolFenceFleetProofToken,
updated_at: OffsetDateTime,
) -> Result<OffsetDateTime> {
update_quota_if_incarnation_stamped(bucket, data, expected_incarnation_id, proof, Some(updated_at)).await
}
async fn update_quota_if_incarnation_stamped(
bucket: &str,
data: Vec<u8>,
expected_incarnation_id: Uuid,
proof: &crate::services::notification_sys::CrossPoolFenceFleetProofToken,
updated_at: Option<OffsetDateTime>,
) -> Result<OffsetDateTime> {
let sys = get_bucket_metadata_sys()?;
let guard = Box::pin(acquire_config_write_guard_for_incarnation(
@@ -821,7 +888,7 @@ pub async fn update_quota_if_incarnation(
achieved: 0,
});
}
update_under_config_write_guard(sys, &guard, rustfs_config::QUOTA_CONFIG_FILE, data).await
update_under_config_write_guard(sys, &guard, rustfs_config::QUOTA_CONFIG_FILE, data, updated_at).await
}
pub async fn update_bucket_targets_under_transaction_lock(
@@ -837,6 +904,7 @@ async fn update_under_config_write_guard(
guard: &BucketMetadataMutationGuard,
config_file: &str,
data: Vec<u8>,
updated_at: Option<OffsetDateTime>,
) -> Result<OffsetDateTime> {
guard.ensure_valid(&guard.bucket)?;
let metadata_sys = sys.read().await.clone();
@@ -848,7 +916,7 @@ async fn update_under_config_write_guard(
Some(&guard.transaction_guard),
&guard.bucket,
"bucket config transaction",
metadata_sys.update_checked(&guard.bucket, config_file, data, true, guard.incarnation_id),
metadata_sys.update_checked(&guard.bucket, config_file, data, true, guard.incarnation_id, updated_at),
),
)
.await?;
@@ -871,7 +939,7 @@ async fn delete_under_config_write_guard(
Some(&guard.transaction_guard),
&guard.bucket,
"bucket config deletion transaction",
metadata_sys.update_checked(&guard.bucket, config_file, Vec::new(), false, guard.incarnation_id),
metadata_sys.update_checked(&guard.bucket, config_file, Vec::new(), false, guard.incarnation_id, None),
),
)
.await?;
@@ -1770,15 +1838,17 @@ impl BucketMetadataSys {
/// `update` and the config read alone). Keep these boxed.
pub async fn update(&self, bucket: &str, config_file: &str, data: Vec<u8>) -> Result<OffsetDateTime> {
let incarnation_id = Box::pin(self.get_bucket_incarnation_id(bucket)).await?;
Box::pin(self.update_checked(bucket, config_file, data, true, incarnation_id)).await
Box::pin(self.update_checked(bucket, config_file, data, true, incarnation_id, None)).await
}
pub async fn delete(&self, bucket: &str, config_file: &str) -> Result<OffsetDateTime> {
let incarnation_id = self.get_bucket_incarnation_id(bucket).await?;
self.update_checked(bucket, config_file, Vec::new(), false, incarnation_id)
self.update_checked(bucket, config_file, Vec::new(), false, incarnation_id, None)
.await
}
/// `updated_at`: `None` stamps the local clock; `Some` persists a
/// replicated edit's source time (backlog#2292).
async fn update_checked(
&self,
bucket: &str,
@@ -1786,6 +1856,7 @@ impl BucketMetadataSys {
data: Vec<u8>,
parse: bool,
expected_incarnation_id: Uuid,
updated_at: Option<OffsetDateTime>,
) -> Result<OffsetDateTime> {
// Load through this system's own store, the one `save` persists to
// (backlog#1052 S7). Reading from the ambient handle instead made the
@@ -1796,7 +1867,10 @@ impl BucketMetadataSys {
return Err(Error::BucketNotFound(bucket.to_string()));
}
let updated = bm.update_config(config_file, data)?;
let updated = match updated_at {
Some(updated_at) => bm.update_config_at(config_file, data, updated_at)?,
None => bm.update_config(config_file, data)?,
};
Box::pin(self.save(bm)).await?;
@@ -3765,6 +3839,57 @@ mod tests {
);
}
/// backlog#2292: the explicit-stamp write path persists the given source
/// time as the config's `*_config_updated_at` — through the incarnation
/// path and through an already-held transaction guard — and survives a
/// reload from disk, while the plain path keeps stamping the local clock.
#[tokio::test]
async fn explicit_updated_at_is_persisted_as_the_config_stamp() {
let (dirs, ecstore) = isolated_store_over_temp_disks().await;
let bucket = "source-stamped-config";
for dir in &dirs {
std::fs::create_dir_all(dir.path().join(bucket)).expect("bucket volume should be created");
}
let sys = Arc::new(RwLock::new(BucketMetadataSys::new(ecstore)));
let source_time = OffsetDateTime::now_utc() - Duration::from_secs(3 * 3600);
let policy = br#"{"Version":"2012-10-17","Statement":[]}"#.to_vec();
let tagging = b"<Tagging><TagSet><Tag><Key>k</Key><Value>v</Value></Tag></TagSet></Tagging>".to_vec();
// Incarnation path (`update_if_incarnation_at` minus the ambient lookup).
let stamped =
update_with_sys_expected(sys.clone(), bucket, BUCKET_POLICY_CONFIG, policy.clone(), None, Some(source_time))
.await
.expect("source-stamped policy write should persist");
assert_eq!(stamped, source_time);
// Held-guard path (`update_under_transaction_lock_at` minus the ambient lookup).
let guard = acquire_config_write_guard(sys.clone(), bucket).await.expect("write guard");
let stamped = update_under_config_write_guard(sys.clone(), &guard, BUCKET_TAGGING_CONFIG, tagging, Some(source_time))
.await
.expect("source-stamped tagging write should persist");
drop(guard);
assert_eq!(stamped, source_time);
let metadata_sys = sys.read().await.clone();
metadata_sys.metadata_map.write().await.clear();
let reloaded = metadata_sys.get_config_from_disk(bucket).await.expect("reload from disk");
assert_eq!(reloaded.policy_config_updated_at, source_time);
assert_eq!(reloaded.tagging_config_updated_at, source_time);
// The plain path is unchanged: a local edit is stamped with the local clock.
let before = OffsetDateTime::now_utc();
let stamped = update_with_sys(sys.clone(), bucket, BUCKET_POLICY_CONFIG, policy)
.await
.expect("locally stamped policy write should persist");
assert!(stamped >= before, "the plain write path must keep stamping the local clock");
let reloaded = metadata_sys.get_config_from_disk(bucket).await.expect("reload from disk");
assert_eq!(reloaded.policy_config_updated_at, stamped);
assert_eq!(
reloaded.tagging_config_updated_at, source_time,
"an unrelated config keeps its source stamp"
);
}
/// The load and the persisted write share one write guard, so concurrent
/// rewrites of the same config compose instead of clobbering each other.
/// Moving the load outside that guard loses all but the last tag.
@@ -3981,10 +4106,16 @@ mod tests {
let new_incarnation = store.bucket_incarnation_id_from_disk(bucket).await.unwrap();
assert_ne!(old_incarnation, new_incarnation);
let err =
update_with_sys_expected(sys.clone(), bucket, BUCKET_TAGGING_CONFIG, b"<Tagging/>".to_vec(), Some(old_incarnation))
.await
.expect_err("a request authorized for the deleted incarnation must fail closed");
let err = update_with_sys_expected(
sys.clone(),
bucket,
BUCKET_TAGGING_CONFIG,
b"<Tagging/>".to_vec(),
Some(old_incarnation),
None,
)
.await
.expect_err("a request authorized for the deleted incarnation must fail closed");
assert!(matches!(err, Error::BucketNotFound(name) if name == bucket));
let persisted = sys.read().await.get_config_from_disk(bucket).await.unwrap();
@@ -4019,7 +4150,7 @@ mod tests {
}],
})
.unwrap();
update_under_config_write_guard(sys, &guard, BUCKET_TAGGING_CONFIG, tagging)
update_under_config_write_guard(sys, &guard, BUCKET_TAGGING_CONFIG, tagging, None)
.await
.unwrap();
assert!(!delete.is_finished());
@@ -20,9 +20,9 @@ pub use rustfs_replication::{
pub(crate) use rustfs_replication::{
ReplicationDeleteSource, ReplicationMultipartPartInput, ReplicationResyncTargetObject, delete_marker_purge_mrf_entry,
delete_marker_purge_version_id, delete_replication_creates_marker, delete_replication_missing_source_decision,
delete_replication_object_opts, heal_uses_delete_replication_path, is_object_lock_denied_delete,
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match,
replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error,
resync_existing_delete_replication_info, resync_target_for_object, should_retry_delete_marker_purge,
single_part_replica_etag_mismatch, target_delete_version_id,
delete_replication_object_opts, delete_replication_target_version_id, heal_uses_delete_replication_path,
is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication,
replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size,
replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info,
resync_target_for_object, should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
};
@@ -882,6 +882,20 @@ fn reconstructed_heal_delete_info(
) -> DeletedObjectReplicationInfo {
let mut rstate = oi.replication_state();
rstate.replicate_decision_str = dsc.to_string();
// The caller hands us a blank ObjectInfo (the source marker may already be
// gone), so the state above carries no target-assigned marker version ids.
// Restore them from the journal: `delete_marker_purge_version_id` must hit
// the id the target reported, not fall back to the source marker id, which
// a target that mints its own ids answers with an idempotent 204 that would
// acknowledge the intent while the real marker stays behind (backlog#2290).
// The corrupt flag rides along so a refusal stays a refusal after restart.
for (arn, version_id) in &entry.target_delete_marker_version_ids {
rstate
.target_delete_marker_version_ids
.entry(arn.clone())
.or_insert_with(|| version_id.clone());
}
rstate.target_delete_marker_version_ids_corrupt |= entry.target_delete_marker_version_ids_corrupt;
let delete_marker_mtime = entry
.delete_marker_mtime
@@ -6601,4 +6615,87 @@ mod tests {
replacement_data
);
}
/// backlog#2290: a delete-marker purge intent that survives a restart
/// through the MRF journal addresses the marker version the TARGET
/// assigned, exactly as the live watcher does (see the
/// `requires_delayed_purge` spawn). The journal carries the per-ARN ids
/// (`targetDeleteMarkerVersionIDs`) and replay restores them into the
/// reconstructed replication state; without that the replay would fall
/// back to the source marker id, which a target that mints its own ids
/// answers with an idempotent 204 — the entry would be acknowledged while
/// the real marker stayed behind.
#[test]
fn mrf_delete_marker_purge_replay_preserves_target_assigned_marker_version() {
use super::super::replication_object_decision_boundary::{delete_marker_purge_mrf_entry, delete_marker_purge_version_id};
let arn = "arn:minio:replication::generic-target:photos".to_string();
let source_marker = uuid::Uuid::new_v4();
let remote_marker = "remote-assigned-marker-version".to_string();
let live_oi = ObjectInfo {
bucket: "photos".to_string(),
name: "obj".to_string(),
version_id: Some(source_marker),
delete_marker: true,
..Default::default()
};
let mut live_state = live_oi.replication_state();
live_state.replicate_decision_str = replicate_decision_for_admitted_targets(std::slice::from_ref(&arn)).to_string();
live_state
.target_delete_marker_version_ids
.insert(arn.clone(), remote_marker.clone());
let live = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: "obj".to_string(),
delete_marker: true,
delete_marker_version_id: Some(source_marker),
replication_state: Some(live_state),
..Default::default()
},
bucket: "photos".to_string(),
..Default::default()
};
assert_eq!(
delete_marker_purge_version_id(live.delete_object.replication_state.as_ref(), &arn, source_marker),
Some(Some(remote_marker.clone())),
"the live purge addresses the recorded target version"
);
// Watch window exhausted: persist the intent, restart, replay it.
let entry = delete_marker_purge_mrf_entry(&live, vec![arn.clone()]);
let replay_oi = ObjectInfo {
bucket: entry.bucket.clone(),
name: entry.object.clone(),
version_id: entry.version_id,
delete_marker: entry.delete_marker,
..Default::default()
};
let dsc = replicate_decision_for_admitted_targets(&entry.target_arns);
let replayed = reconstructed_heal_delete_info(&entry, &replay_oi, &dsc);
assert_eq!(
delete_marker_purge_version_id(replayed.delete_object.replication_state.as_ref(), &arn, source_marker),
Some(Some(remote_marker)),
"the MRF replay must address the target-assigned marker version, not source marker {source_marker}"
);
// A refusal (inconsistent recorded ids) must stay a refusal across the
// journal round trip instead of degrading into the source-id fallback.
let mut refused = live;
refused
.delete_object
.replication_state
.as_mut()
.expect("state was set above")
.target_delete_marker_version_ids_corrupt = true;
let entry = delete_marker_purge_mrf_entry(&refused, vec![arn.clone()]);
assert!(entry.target_delete_marker_version_ids_corrupt);
let replayed = reconstructed_heal_delete_info(&entry, &replay_oi, &dsc);
assert_eq!(
delete_marker_purge_version_id(replayed.delete_object.replication_state.as_ref(), &arn, source_marker),
None,
"the MRF replay must keep refusing to guess when the recorded ids were inconsistent"
);
}
}
@@ -32,11 +32,11 @@ use super::replication_msgp_boundary::ReplicationMsgpCodec;
use super::replication_object_config::{ReplicationConfig, get_replication_config, must_replicate};
use super::replication_object_decision_boundary::{
MustReplicateOptions, ReplicationMultipartPartInput, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
delete_replication_creates_marker, heal_uses_delete_replication_path, is_object_lock_denied_delete,
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match,
replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error,
resync_existing_delete_replication_info, should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
target_delete_version_id,
delete_replication_creates_marker, delete_replication_target_version_id, heal_uses_delete_replication_path,
is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication,
replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size,
replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info,
should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
};
use super::replication_queue_boundary::{DeletedObjectReplicationInfo, ReplicationQueueAdmission};
use super::replication_resync_boundary::ResyncStatusType;
@@ -2051,7 +2051,11 @@ pub(crate) async fn replicate_delete_with_outcome<S: ReplicationStorage>(
let is_version_purge = is_version_delete_replication(&dobj.delete_object);
let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object);
// The watcher exists to purge a replicated marker once the SOURCE marker
// vanishes. A version purge is that purge already (its failures reach the
// journal as a purge entry), so it must not spawn a second watcher that
// journals a duplicate intent (backlog#2290).
let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object) && !is_version_purge;
let (replication_status, prev_status) = if !is_version_purge {
(
@@ -2761,12 +2765,6 @@ fn unavailable_delete_target_info(dobj: &DeletedObjectReplicationInfo, arn: &str
}
async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_client: Arc<TargetClient>) -> ReplicatedTargetInfo {
let version_id = if let Some(version_id) = &dobj.delete_object.delete_marker_version_id {
version_id.to_owned()
} else {
dobj.delete_object.version_id.unwrap_or_default()
};
let mut rinfo = dobj
.delete_object
.replication_state
@@ -2799,7 +2797,25 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
return rinfo;
}
let version_id = target_delete_version_id(version_id, is_version_purge);
// Purging a replicated delete marker addresses the version the target
// assigned (recorded when the marker was created there); see
// `delete_replication_target_version_id`. A corrupt record is a failure,
// not a guess: the entry stays visible until the metadata is repaired.
let Some(version_id) = delete_replication_target_version_id(&dobj.delete_object, &tgt_client.arn) else {
warn!(
event = EVENT_DELETE_MARKER_PURGE_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = tgt_client.bucket,
object = dobj.delete_object.object_name,
arn = %tgt_client.arn,
reason = "recorded_target_version_inconsistent",
"Replicated version purge refused: recorded target delete-marker version metadata is inconsistent"
);
rinfo.version_purge_status = VersionPurgeStatusType::Failed;
rinfo.error = Some("recorded target delete-marker version metadata is inconsistent".to_string());
return rinfo;
};
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match head_object_for_worker(
+34 -8
View File
@@ -429,6 +429,27 @@ where
}
}
/// The cached mapping record for one user or group, looked up in the same
/// cache partition `policy_db_set` writes it to (group / STS / regular+service
/// user). `None` when no mapping is stored.
pub async fn get_mapped_policy_record(&self, name: &str, user_type: UserType, is_group: bool) -> Option<MappedPolicy> {
let cache = self.cache.snapshot();
if is_group {
cache.group_policies.get(name).cloned()
} else if user_type == UserType::Sts {
cache.sts_policies.get(name).cloned()
} else {
cache.user_policies.get(name).cloned()
}
}
/// The cached group record (members, status, own timestamp) without the
/// mapped-policy overlay `get_group_description` applies. `None` when the
/// group does not exist.
pub async fn get_group_info(&self, name: &str) -> Option<GroupInfo> {
self.cache.snapshot().groups.get(name).cloned()
}
pub async fn get_policy(&self, name: &str) -> Result<Policy> {
if name.is_empty() {
return Err(Error::InvalidArgument);
@@ -1693,6 +1714,10 @@ where
}
}
// The group's own timestamp moves with every membership or status
// change: site replication judges an incoming group item against it
// (backlog#2291), so it must reflect the last change, not creation.
let now = OffsetDateTime::now_utc();
let gi = match cache.groups.get(group) {
Some(res) => {
let mut gi = res.clone();
@@ -1701,6 +1726,7 @@ where
uniq_set.extend(members.iter().cloned());
gi.members = uniq_set.into_iter().collect();
gi.update_at = Some(now);
gi
}
None => GroupInfo::new(members.clone()),
@@ -1709,8 +1735,7 @@ where
self.api.save_group_info(group, gi.clone()).await?;
let now = self.cache.with_write_lock(|cache| {
let now = OffsetDateTime::now_utc();
self.cache.with_write_lock(|cache| {
cache.add_or_update_group(group, &gi, now);
let user_group_memberships = Arc::clone(&cache.state().user_group_memberships);
@@ -1719,7 +1744,6 @@ where
m.insert(group.to_string());
cache.add_or_update_user_group_membership(member, &m, now);
});
now
});
Ok(now)
@@ -1743,12 +1767,14 @@ where
} else {
gi.status = STATUS_DISABLED.to_owned();
}
let now = OffsetDateTime::now_utc();
gi.update_at = Some(now);
self.api.save_group_info(name, gi.clone()).await?;
self.cache.add_or_update_group(name, &gi, OffsetDateTime::now_utc());
self.cache.add_or_update_group(name, &gi, now);
Ok(OffsetDateTime::now_utc())
Ok(now)
}
pub async fn get_group_description(&self, name: &str) -> Result<GroupDesc> {
@@ -1830,13 +1856,14 @@ where
let s: HashSet<&String> = HashSet::from_iter(gi.members.iter());
let d: HashSet<&String> = HashSet::from_iter(members.iter());
gi.members = s.difference(&d).map(|v| v.to_string()).collect::<Vec<String>>();
let now = OffsetDateTime::now_utc();
gi.update_at = Some(now);
if !update_cache_only {
self.api.save_group_info(name, gi.clone()).await?;
}
let now = self.cache.with_write_lock(|cache| {
let now = OffsetDateTime::now_utc();
self.cache.with_write_lock(|cache| {
cache.add_or_update_group(name, &gi, now);
let user_group_memberships = Arc::clone(&cache.state().user_group_memberships);
@@ -1847,7 +1874,6 @@ where
cache.add_or_update_user_group_membership(member, &m, now);
}
});
now
});
Ok(now)
+16
View File
@@ -1055,6 +1055,22 @@ impl<T: Store> IamSys<T> {
self.store.get_group_description(group).await
}
/// The stored group record itself (see `IamCache::get_group_info`).
pub async fn get_group_info(&self, group: &str) -> Option<GroupInfo> {
self.store.get_group_info(group).await
}
/// The stored policy document, `Error::NoSuchPolicy` when absent.
pub async fn get_policy_doc(&self, name: &str) -> Result<PolicyDoc> {
self.store.get_policy_doc(name).await
}
/// The stored mapping record for one user or group (see
/// `IamCache::get_mapped_policy_record`).
pub async fn get_mapped_policy_record(&self, name: &str, user_type: UserType, is_group: bool) -> Option<MappedPolicy> {
self.store.get_mapped_policy_record(name, user_type, is_group).await
}
pub async fn list_groups_load(&self) -> Result<Vec<String>> {
self.store.update_groups().await
}
+163 -3
View File
@@ -76,6 +76,21 @@ impl ReplicationWorkerOperation for DeletedObjectReplicationInfo {
.delete_object
.delete_marker_mtime
.and_then(|t| i64::try_from(t.unix_timestamp_nanos()).ok()),
// Carry the target-assigned marker version ids (and the fail-closed corrupt
// flag) into the journal so a purge intent replayed after a restart addresses
// the same version the live path did (backlog#2290). Only delete-marker state
// ever records these; other deletes serialize an empty map.
target_delete_marker_version_ids: self
.delete_object
.replication_state
.as_ref()
.map(|state| state.target_delete_marker_version_ids.clone())
.unwrap_or_default(),
target_delete_marker_version_ids_corrupt: self
.delete_object
.replication_state
.as_ref()
.is_some_and(|state| state.target_delete_marker_version_ids_corrupt),
target_arns: self.admitted_target_arns(),
force_delete_id: self.delete_object.force_delete_id,
force_delete_generation: self.delete_object.force_delete_generation,
@@ -238,6 +253,28 @@ pub fn delete_marker_purge_version_id(
})
}
/// The version a delete replication addresses on `arn`, or `None` to refuse.
///
/// A version purge whose purged version is a delete marker must address the
/// marker version the TARGET assigned — the recorded mapping, exactly as the
/// delayed-purge watcher does. The source-side `DELETE ?versionId=<marker>`
/// replicates as such a purge, and a generic S3 target answers a DELETE of an
/// unknown versionId with 204 while keeping its marker, so addressing it by
/// the source id reported success and left the marker behind (backlog#2290,
/// R6.1 on the VMs). Nothing recorded falls back to the source-derived id
/// (id-mirroring peers); a corrupt record refuses, as the watcher does.
pub fn delete_replication_target_version_id(dobj: &DeletedObject, arn: &str) -> Option<Option<String>> {
let is_version_purge = is_version_delete_replication(dobj);
if is_version_purge
&& !dobj.delete_marker
&& let Some(marker) = dobj.delete_marker_version_id
{
return delete_marker_purge_version_id(dobj.replication_state.as_ref(), arn, marker);
}
let source_version = dobj.delete_marker_version_id.or(dobj.version_id).unwrap_or_default();
Some(target_delete_version_id(source_version, is_version_purge))
}
/// Shape an exhausted purge intent as a marker-creation delete entry. Replay
/// reconstructs it with `delete_marker: true`, finds the source marker gone,
/// and funnels into the stale-marker branch of `replicate_delete_with_outcome`
@@ -258,9 +295,9 @@ mod tests {
use super::{
DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error,
is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info,
should_retry_delete_marker_purge, target_delete_version_id,
delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete,
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome,
resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id,
};
use crate::storage_api::DeletedObject;
use crate::{
@@ -595,6 +632,76 @@ mod tests {
assert_eq!(entry.retry_count, 0);
assert_eq!(entry.bucket, "bucket-a");
assert_eq!(entry.object, "doc.txt");
assert!(
entry.target_delete_marker_version_ids.is_empty(),
"no recorded target marker ids means the journal carries none"
);
assert!(!entry.target_delete_marker_version_ids_corrupt);
}
/// backlog#2290: a purge intent journaled to MRF must carry the marker
/// version ids the targets assigned, plus the fail-closed corrupt flag,
/// so a replay after restart addresses the same version the live path did.
#[test]
fn delete_marker_purge_mrf_entry_carries_target_assigned_marker_versions() {
let delete_marker_version_id = Uuid::new_v4();
let mut state = ReplicationState::default();
state
.target_delete_marker_version_ids
.insert("arn:a".to_string(), "remote-marker-a".to_string());
state
.target_delete_marker_version_ids
.insert("arn:b".to_string(), "remote-marker-b".to_string());
let mut dobj = DeletedObjectReplicationInfo {
delete_object: DeletedObject {
object_name: "doc.txt".to_string(),
delete_marker: false,
version_id: Some(Uuid::new_v4()),
delete_marker_version_id: Some(delete_marker_version_id),
replication_state: Some(state),
..Default::default()
},
bucket: "bucket-a".to_string(),
..Default::default()
};
let entry = delete_marker_purge_mrf_entry(&dobj, vec!["arn:a".to_string()]);
assert_eq!(
entry.target_delete_marker_version_ids,
HashMap::from([
("arn:a".to_string(), "remote-marker-a".to_string()),
("arn:b".to_string(), "remote-marker-b".to_string()),
]),
"every recorded target marker id survives the journal, regardless of the retried ARN subset"
);
assert!(!entry.target_delete_marker_version_ids_corrupt);
assert_eq!(
delete_marker_purge_version_id(
Some(&ReplicationState {
target_delete_marker_version_ids: entry.target_delete_marker_version_ids,
..Default::default()
}),
"arn:a",
delete_marker_version_id
),
Some(Some("remote-marker-a".to_string()))
);
// The live path refuses to purge on inconsistent metadata and reports the target
// as failed; the journaled intent must keep refusing after a restart.
dobj.delete_object
.replication_state
.as_mut()
.expect("state was set above")
.target_delete_marker_version_ids_corrupt = true;
let entry = delete_marker_purge_mrf_entry(&dobj, vec!["arn:a".to_string()]);
assert!(entry.target_delete_marker_version_ids_corrupt);
// A delete without replication state journals an empty map.
dobj.delete_object.replication_state = None;
let entry = dobj.to_mrf_entry();
assert!(entry.target_delete_marker_version_ids.is_empty());
assert!(!entry.target_delete_marker_version_ids_corrupt);
}
#[test]
@@ -656,4 +763,57 @@ mod tests {
assert!(!is_object_lock_denied_delete(Some("InternalError"), Some("retention lookup failed")));
assert!(!is_object_lock_denied_delete(None, Some("legal hold")));
}
fn purge_of_marker(marker: Uuid, state: Option<ReplicationState>) -> DeletedObject {
DeletedObject {
object_name: "obj".to_string(),
delete_marker: false,
delete_marker_version_id: Some(marker),
version_id: None,
replication_state: state,
..Default::default()
}
}
#[test]
fn delete_replication_target_version_id_addresses_recorded_marker_for_purges() {
let arn = "arn:minio:replication::generic:photos";
let marker = Uuid::new_v4();
let mut state = ReplicationState::default();
state
.target_delete_marker_version_ids
.insert(arn.to_string(), "remote-marker".to_string());
// purge of a replicated marker: the target's own version
assert_eq!(
delete_replication_target_version_id(&purge_of_marker(marker, Some(state.clone())), arn),
Some(Some("remote-marker".to_string()))
);
// nothing recorded for this arn: the source-derived id (id-mirroring peers)
assert_eq!(
delete_replication_target_version_id(&purge_of_marker(marker, None), arn),
Some(Some(marker.to_string()))
);
// corrupt record: refuse instead of guessing
state.target_delete_marker_version_ids_corrupt = true;
assert_eq!(delete_replication_target_version_id(&purge_of_marker(marker, Some(state)), arn), None);
// marker creation keeps the source id (the target mints its own on a
// versionless DELETE; the id only travels in the source header)
let creation = DeletedObject {
object_name: "obj".to_string(),
delete_marker: true,
delete_marker_version_id: Some(marker),
..Default::default()
};
assert_eq!(delete_replication_target_version_id(&creation, arn), Some(Some(marker.to_string())));
// plain version purge: the source version id
let version = Uuid::new_v4();
let purge = DeletedObject {
object_name: "obj".to_string(),
version_id: Some(version),
..Default::default()
};
assert_eq!(delete_replication_target_version_id(&purge, arn), Some(Some(version.to_string())));
}
}
+20
View File
@@ -641,6 +641,26 @@ pub struct MrfReplicateEntry {
#[serde(rename = "deleteMarkerMtime", skip_serializing_if = "Option::is_none", default)]
pub delete_marker_mtime: Option<i64>,
// For delete-marker purge intents: the exact version id each target assigned to the
// replicated marker, keyed by target ARN. A generic S3 target mints its own version ids
// and answers a DELETE of an unknown id with 204, so a replay that fell back to the source
// marker id would be acknowledged while the real marker stayed behind (backlog#2290).
// Old files lack this key; default=empty means "unknown" and replay keeps the source-id
// fallback it always had.
#[serde(rename = "targetDeleteMarkerVersionIDs", skip_serializing_if = "HashMap::is_empty", default)]
pub target_delete_marker_version_ids: HashMap<String, String>,
// Companion to the map above: the source metadata disagreed about the recorded ids when
// the intent was journaled, so the live path refused to guess and reported the target as
// failed. Replay must keep refusing instead of falling back to the source id. Old files
// lack this key; default=false.
#[serde(
rename = "targetDeleteMarkerVersionIDsCorrupt",
skip_serializing_if = "std::ops::Not::not",
default
)]
pub target_delete_marker_version_ids_corrupt: bool,
#[serde(rename = "targetARNs", skip_serializing_if = "Vec::is_empty", default)]
pub target_arns: Vec<String>,
+3 -3
View File
@@ -41,9 +41,9 @@ pub use config::{
};
pub use delete::{
DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error,
is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info,
should_retry_delete_marker_purge, target_delete_version_id,
delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete,
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome,
resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id,
};
pub use filemeta::{
NULL_VERSION_ID, REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING,
+167 -2
View File
@@ -31,8 +31,13 @@ const CAPABILITY_OPERATION_KIND: u64 = 1 << 0;
const CAPABILITY_TARGET_ARNS: u64 = 1 << 1;
const CAPABILITY_FORCE_DELETE: u64 = 1 << 2;
const CAPABILITY_DELETE_MARKER_MTIME: u64 = 1 << 3;
const MRF_KNOWN_CAPABILITIES: u64 =
CAPABILITY_OPERATION_KIND | CAPABILITY_TARGET_ARNS | CAPABILITY_FORCE_DELETE | CAPABILITY_DELETE_MARKER_MTIME;
// Per-ARN target-assigned delete-marker version ids on purge intents (backlog#2290).
const CAPABILITY_TARGET_DELETE_MARKER_VERSION_IDS: u64 = 1 << 4;
const MRF_KNOWN_CAPABILITIES: u64 = CAPABILITY_OPERATION_KIND
| CAPABILITY_TARGET_ARNS
| CAPABILITY_FORCE_DELETE
| CAPABILITY_DELETE_MARKER_MTIME
| CAPABILITY_TARGET_DELETE_MARKER_VERSION_IDS;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MrfCapability {
@@ -40,6 +45,7 @@ pub enum MrfCapability {
TargetArns,
ForceDelete,
DeleteMarkerMtime,
TargetDeleteMarkerVersionIds,
}
impl MrfCapability {
@@ -49,6 +55,7 @@ impl MrfCapability {
Self::TargetArns => CAPABILITY_TARGET_ARNS,
Self::ForceDelete => CAPABILITY_FORCE_DELETE,
Self::DeleteMarkerMtime => CAPABILITY_DELETE_MARKER_MTIME,
Self::TargetDeleteMarkerVersionIds => CAPABILITY_TARGET_DELETE_MARKER_VERSION_IDS,
}
}
}
@@ -601,9 +608,17 @@ pub fn decode_mrf_file(data: &[u8]) -> Result<Vec<MrfReplicateEntry>> {
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use uuid::Uuid;
// Capability word 31 = OperationKind | TargetArns | ForceDelete | DeleteMarkerMtime |
// TargetDeleteMarkerVersionIds (backlog#2290).
const ENVELOPE_FIXTURE: &[u8] = &[
b'M', b'R', b'F', b'E', 1, 0, 1, 0, 1, 0, 0, 0, 31, 0, 0, 0, 0, 0, 0, 0, 3, 0, 0, 0, 1, 2, 3,
];
// The envelope a binary from before backlog#2290 writes: same header, capability word 15.
const PRE_TARGET_MARKER_IDS_ENVELOPE_FIXTURE: &[u8] = &[
b'M', b'R', b'F', b'E', 1, 0, 1, 0, 1, 0, 0, 0, 15, 0, 0, 0, 0, 0, 0, 0, 3, 0, 0, 0, 1, 2, 3,
];
@@ -626,6 +641,8 @@ mod tests {
delete_marker_version_id: None,
delete_marker: false,
delete_marker_mtime: None,
target_delete_marker_version_ids: HashMap::new(),
target_delete_marker_version_ids_corrupt: false,
target_arns: vec!["arn:target-a".to_string()],
},
MrfReplicateEntry {
@@ -642,6 +659,8 @@ mod tests {
delete_marker_version_id: None,
delete_marker: false,
delete_marker_mtime: None,
target_delete_marker_version_ids: HashMap::new(),
target_delete_marker_version_ids_corrupt: false,
target_arns: vec!["arn:target-a".to_string(), "arn:target-b".to_string()],
},
MrfReplicateEntry {
@@ -658,6 +677,11 @@ mod tests {
delete_marker_version_id: Some(del_vid),
delete_marker: true,
delete_marker_mtime: Some(1_705_312_200_123_456_789),
target_delete_marker_version_ids: HashMap::from([
("arn:target-a".to_string(), "remote-marker-a".to_string()),
("arn:target-b".to_string(), "remote-marker-b".to_string()),
]),
target_delete_marker_version_ids_corrupt: false,
target_arns: vec!["arn:target-a".to_string()],
},
];
@@ -685,6 +709,54 @@ mod tests {
Some(1_705_312_200_123_456_789),
"delete-marker mtime must survive the MRF disk round-trip"
);
assert!(decoded[0].target_delete_marker_version_ids.is_empty());
assert!(decoded[1].target_delete_marker_version_ids.is_empty());
assert_eq!(
decoded[2].target_delete_marker_version_ids,
HashMap::from([
("arn:target-a".to_string(), "remote-marker-a".to_string()),
("arn:target-b".to_string(), "remote-marker-b".to_string()),
]),
"target-assigned marker version ids must survive the MRF disk round-trip (backlog#2290)"
);
assert!(!decoded[2].target_delete_marker_version_ids_corrupt);
}
/// backlog#2290: the corrupt flag rides the same journal round trip, and an
/// entry that carries neither field encodes exactly as it did before the
/// field existed (both keys are skipped when empty/false).
#[test]
fn mrf_file_round_trips_target_marker_ids_corrupt_flag_and_skips_empty_keys() {
let corrupt = MrfReplicateEntry {
bucket: "bucket-a".to_string(),
object: "delete-a".to_string(),
op: MrfOpKind::Delete,
delete_marker: true,
delete_marker_version_id: Some(Uuid::new_v4()),
target_delete_marker_version_ids_corrupt: true,
target_arns: vec!["arn:target-a".to_string()],
..Default::default()
};
let decoded = decode_mrf_file(&encode_mrf_file(std::slice::from_ref(&corrupt)).expect("mrf file should encode"))
.expect("mrf file should decode");
assert_eq!(decoded, vec![corrupt]);
assert!(decoded[0].target_delete_marker_version_ids_corrupt);
let plain = MrfReplicateEntry {
bucket: "bucket-a".to_string(),
object: "delete-a".to_string(),
op: MrfOpKind::Delete,
delete_marker: true,
target_arns: vec!["arn:target-a".to_string()],
..Default::default()
};
let encoded = encode_mrf_file(std::slice::from_ref(&plain)).expect("mrf file should encode");
let payload = String::from_utf8_lossy(&encoded);
assert!(
!payload.contains("targetDeleteMarkerVersionIDs"),
"an entry without recorded ids must not grow the new keys: {payload}"
);
assert_eq!(decode_mrf_file(&encoded).expect("mrf file should decode"), vec![plain]);
}
#[test]
@@ -719,6 +791,99 @@ mod tests {
// Old files lack the deleteMarkerMtime key; it must default to None so replay keeps the
// pre-#867 fallback to the current time.
assert_eq!(decoded[0].delete_marker_mtime, None);
// Old files also lack the target marker id keys; they must default to an empty map
// and a clear corrupt flag so replay keeps the pre-#2290 source-id fallback.
assert!(decoded[0].target_delete_marker_version_ids.is_empty());
assert!(!decoded[0].target_delete_marker_version_ids_corrupt);
}
/// backlog#2290: a delete-marker entry written by a binary that predates the
/// `targetDeleteMarkerVersionIDs` key decodes with an empty map and a clear
/// corrupt flag — the exact shape replay handled before the field existed.
#[test]
fn mrf_pre_target_marker_ids_delete_entry_decodes_with_empty_map() {
let marker_version_id = Uuid::new_v4();
let mut payload = Vec::new();
rmp::encode::write_array_len(&mut payload, 1).expect("array len should encode");
rmp::encode::write_map_len(&mut payload, 9).expect("map len should encode");
rmp::encode::write_str(&mut payload, "bucket").expect("bucket key should encode");
rmp::encode::write_str(&mut payload, "old-bucket").expect("bucket value should encode");
rmp::encode::write_str(&mut payload, "object").expect("object key should encode");
rmp::encode::write_str(&mut payload, "old-key").expect("object value should encode");
rmp::encode::write_str(&mut payload, "retryCount").expect("retry key should encode");
rmp::encode::write_i32(&mut payload, 0).expect("retry value should encode");
rmp::encode::write_str(&mut payload, "size").expect("size key should encode");
rmp::encode::write_i64(&mut payload, 0).expect("size value should encode");
rmp::encode::write_str(&mut payload, "op").expect("op key should encode");
rmp::encode::write_str(&mut payload, "delete").expect("op value should encode");
rmp::encode::write_str(&mut payload, "forceDelete").expect("forceDelete key should encode");
rmp::encode::write_bool(&mut payload, false).expect("forceDelete value should encode");
rmp::encode::write_str(&mut payload, "deleteMarkerVersionID").expect("marker id key should encode");
// Uuid serializes as a 16-byte bin in the MessagePack journal.
rmp::encode::write_bin(&mut payload, marker_version_id.as_bytes()).expect("marker id value should encode");
rmp::encode::write_str(&mut payload, "deleteMarker").expect("deleteMarker key should encode");
rmp::encode::write_bool(&mut payload, true).expect("deleteMarker value should encode");
rmp::encode::write_str(&mut payload, "targetARNs").expect("targetARNs key should encode");
rmp::encode::write_array_len(&mut payload, 1).expect("targetARNs len should encode");
rmp::encode::write_str(&mut payload, "arn:target-a").expect("targetARNs value should encode");
let mut data = Vec::with_capacity(4 + payload.len());
data.extend_from_slice(&MRF_META_FORMAT.to_le_bytes());
data.extend_from_slice(&MRF_META_VERSION.to_le_bytes());
data.extend_from_slice(&payload);
let decoded = decode_mrf_file(&data).expect("pre-#2290 delete-marker entry should decode");
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].op, MrfOpKind::Delete);
assert!(decoded[0].delete_marker);
assert_eq!(decoded[0].delete_marker_version_id, Some(marker_version_id));
assert_eq!(decoded[0].target_arns, vec!["arn:target-a".to_string()]);
assert!(decoded[0].target_delete_marker_version_ids.is_empty());
assert!(!decoded[0].target_delete_marker_version_ids_corrupt);
}
/// backlog#2290: the new field is fenced by its own capability bit exactly
/// like the earlier optional fields — a reader without the bit refuses an
/// envelope that advertises it, while the current reader still accepts the
/// pre-#2290 envelope.
#[test]
fn envelope_target_marker_ids_capability_is_fenced_and_backward_compatible() {
assert!(MrfCapabilities::current().contains(MrfCapability::TargetDeleteMarkerVersionIds));
assert_eq!(MrfCapabilities::with(MrfCapability::TargetDeleteMarkerVersionIds).bits(), 1 << 4);
// Old envelope, current reader: accepted, and the negotiated set lacks the new bit.
let legacy = MrfEnvelope::decode(PRE_TARGET_MARKER_IDS_ENVELOPE_FIXTURE, MrfProtocolCapabilities::current())
.expect("pre-#2290 envelope should decode");
assert_eq!(legacy.protocol().capabilities().bits(), 15);
assert!(
!legacy
.protocol()
.capabilities()
.contains(MrfCapability::TargetDeleteMarkerVersionIds)
);
assert_eq!(legacy.payload(), &[1, 2, 3]);
// Current envelope, reader that only knows the pre-#2290 bits: refused.
let pre_2290_reader = MrfProtocolCapabilities::new(1, 1, MrfCapabilities::from_bits(15).expect("known bits"));
assert_eq!(
MrfEnvelope::decode(ENVELOPE_FIXTURE, pre_2290_reader),
Err(MrfEnvelopeError::MissingCapabilities {
required: 31,
available: 15,
})
);
// Negotiation with such a peer drops the bit instead of failing.
let negotiated = MrfProtocolCapabilities::current()
.negotiate(pre_2290_reader)
.expect("negotiation with a pre-#2290 peer should succeed");
assert!(
!negotiated
.capabilities()
.contains(MrfCapability::TargetDeleteMarkerVersionIds)
);
assert!(negotiated.capabilities().contains(MrfCapability::DeleteMarkerMtime));
}
#[test]
+627 -80
View File
@@ -66,8 +66,8 @@ use rustfs_madmin::{
ReplicateRemoveStatus, ResyncBucketStatus, SITE_REPL_API_VERSION, SR_IAM_ITEM_STS_ACC, SR_IAM_ITEM_STS_ACC_LEGACY,
SRBucketInfo, SRBucketMeta, SRBucketStatsSummary, SRGroupInfo, SRGroupStatsSummary, SRIAMItem, SRIAMUser,
SRILMExpiryStatsSummary, SRInfo, SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq, SRPendingOperation, SRPolicyMapping,
SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRSTSCredential, SRSessionPolicy, SRSiteSummary, SRStateEditReq,
SRStateInfo, SRStatusInfo, SRSvcAccChange, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRSTSCredential, SRSiteSummary, SRStateEditReq, SRStateInfo,
SRStatusInfo, SRSvcAccChange, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
};
use rustfs_policy::policy::{
Policy,
@@ -87,7 +87,7 @@ use std::sync::{LazyLock, Mutex as StdMutex};
use std::time::Duration;
use time::OffsetDateTime;
use tokio::sync::Mutex;
use tracing::{info, warn};
use tracing::{debug, info, warn};
use url::Url;
use url::form_urlencoded;
use uuid::Uuid;
@@ -98,7 +98,6 @@ use uuid::Uuid;
// paths keep resolving while this file keeps only the HTTP handlers.
pub(crate) use crate::site_replication::*;
const SERVICE_ACCOUNT_ENVELOPE_VERSION: u64 = 2;
// Serializes peer-join admission (staleness check -> IAM upsert -> state
// commit) across every node of this site; see admit_peer_join. Never an
// actual object — only a namespace-lock key, like the repair execution lock.
@@ -1980,7 +1979,7 @@ async fn bootstrap_existing_metadata_after_add(
return errors;
}
};
let plan = match site_replication_bootstrap_plan(&info) {
let plan = match build_site_replication_bootstrap_plan(&info).await {
Ok(plan) => plan,
Err(err) => {
let mut errors = SiteReplicationErrorSummary::default();
@@ -3498,6 +3497,7 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic
state.resync_status.clear();
state.retry_queue.clear();
state.iam_deletion_replays.clear();
state.iam_deletion_marks.clear();
state.pending_endpoint_refresh = None;
state.updated_at = Some(OffsetDateTime::now_utc());
return state;
@@ -3509,6 +3509,7 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic
state.resync_status.clear();
state.retry_queue.clear();
state.iam_deletion_replays.clear();
state.iam_deletion_marks.clear();
state.pending_endpoint_refresh = None;
state.updated_at = Some(OffsetDateTime::now_utc());
return state;
@@ -5386,6 +5387,38 @@ fn is_stale_update(local_updated_at: OffsetDateTime, incoming_updated_at: Option
incoming_updated_at.is_some_and(|incoming_updated_at| incoming_updated_at < local_updated_at)
}
/// Verdict for an incoming IAM item judged against the local record it would
/// overwrite or delete (backlog#2291).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum IamItemVerdict {
/// Apply the item: there is no local record, the item carries no source
/// timestamp (older peer), or it is at least as new as the local record.
Apply,
/// The local record was written from a newer source change; acknowledge the
/// item without touching the record. Covers both directions: a delayed
/// grant must not undo a newer revoke, and a delayed revoke must not undo a
/// newer grant.
SkipStale,
}
/// Ordering rule shared by the `policy`, `policy-mapping` and `group-info`
/// item paths (and matching `iam-user` / `service-account`).
///
/// `local_record_updated_at` is `None` when the targeted record does not
/// exist locally: nothing can be stale relative to an absent record, so a
/// create is applied and a delete falls through to the idempotent no-op paths
/// (backlog#2071). A record that exists but predates timestamps passes
/// `Some(UNIX_EPOCH)` and therefore never rejects an item.
fn judge_iam_item_staleness(
local_record_updated_at: Option<OffsetDateTime>,
incoming_updated_at: Option<OffsetDateTime>,
) -> IamItemVerdict {
match local_record_updated_at {
Some(local_updated_at) if is_stale_update(local_updated_at, incoming_updated_at) => IamItemVerdict::SkipStale,
_ => IamItemVerdict::Apply,
}
}
fn bucket_meta_local_updated_at(
bucket_meta: &crate::admin::storage_api::bucket::metadata::BucketMetadata,
config_file: &str,
@@ -5586,6 +5619,17 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
_ => unreachable!(),
};
// Persist the SOURCE `updated_at` as the stored `*_config_updated_at`
// stamp (backlog#2292). The staleness gate above compares the next item's
// source time against that stamp, so stamping the local apply time would
// reject a newer source edit that was merely delivered after this write
// (two quick edits under delivery delay, or a peer clock ahead of ours).
// Items without a source time keep the local stamp; lc-config keeps it
// too: its staleness axis is the in-document `expiry_updated_at` the merge
// above records, and the whole-config time is only its deletion / legacy
// lower bound.
let source_updated_at = if item.r#type == "lc-config" { None } else { item.updated_at };
if !skip_config_write {
if let Some(data) = data {
if item.r#type == "quota-config" {
@@ -5604,13 +5648,25 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
"durable quota capability is not confirmed across the cluster".to_string(),
)
})?;
metadata_sys::update_quota_if_incarnation(&item.bucket, data, expected_incarnation_id, &proof)
.await
.map_err(ApiError::from)?;
match source_updated_at {
Some(source_updated_at) => {
metadata_sys::update_quota_if_incarnation_at(
&item.bucket,
data,
expected_incarnation_id,
&proof,
source_updated_at,
)
.await
}
None => {
metadata_sys::update_quota_if_incarnation(&item.bucket, data, expected_incarnation_id, &proof).await
}
}
.map_err(ApiError::from)?;
} else {
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
write_replicated_bucket_config(&item.bucket, config_file, data, expected_incarnation_id, source_updated_at)
.await?;
}
} else {
if let Some(guard) = lifecycle_guard.as_ref() {
@@ -5618,9 +5674,8 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
.await
.map_err(ApiError::from)?;
} else {
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
write_replicated_bucket_config(&item.bucket, config_file, data, expected_incarnation_id, source_updated_at)
.await?;
}
}
} else {
@@ -5656,43 +5711,28 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
Ok(())
}
fn group_info_requires_upsert(update: &rustfs_madmin::GroupAddRemove) -> bool {
!update.is_remove
/// Write one replicated bucket config, stamped with the item's source
/// `updated_at` when it carries one and with the local clock otherwise
/// (backlog#2292; see [`apply_bucket_meta_item`]).
async fn write_replicated_bucket_config(
bucket: &str,
config_file: &str,
data: Vec<u8>,
expected_incarnation_id: Uuid,
source_updated_at: Option<OffsetDateTime>,
) -> S3Result<()> {
match source_updated_at {
Some(source_updated_at) => {
metadata_sys::update_if_incarnation_at(bucket, config_file, data, expected_incarnation_id, source_updated_at).await
}
None => metadata_sys::update_if_incarnation(bucket, config_file, data, expected_incarnation_id).await,
}
.map_err(ApiError::from)?;
Ok(())
}
pub(crate) fn encode_service_account_replication_policy(
claims: &HashMap<String, Value>,
session_policy: Option<&str>,
) -> S3Result<(SRSessionPolicy, Option<rustfs_madmin::SRSvcAccReplicationEnvelope>)> {
if !claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM) {
return session_policy
.map(SRSessionPolicy::from_json)
.transpose()
.map(|policy| policy.unwrap_or_default())
.map(|policy| (policy, None))
.map_err(|err| s3_error!(InvalidArgument, "marshal policy failed: {:?}", err));
}
let policy = match session_policy {
Some(policy) => serde_json::from_str::<Policy>(policy)
.map_err(|err| s3_error!(InvalidArgument, "invalid service account replication policy: {:?}", err))?,
None => Policy::default(),
};
if policy.statements.is_empty() && (!policy.id.is_empty() || !policy.version.is_empty())
|| policy.version.is_empty() && !policy.statements.is_empty()
{
return Err(s3_error!(InvalidArgument, "service account replication policy is not normalized"));
}
let policy = serde_json::to_string(&policy)
.map_err(|err| s3_error!(InternalError, "marshal service account replication policy failed: {:?}", err))?;
let policy = SRSessionPolicy::from_json(&policy)
.map_err(|err| s3_error!(InternalError, "marshal service account replication policy failed: {:?}", err))?;
Ok((
policy,
Some(rustfs_madmin::SRSvcAccReplicationEnvelope {
version: SERVICE_ACCOUNT_ENVELOPE_VERSION,
}),
))
fn group_info_requires_upsert(update: &rustfs_madmin::GroupAddRemove) -> bool {
!update.is_remove
}
#[derive(Debug)]
@@ -5762,27 +5802,87 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> {
return Err(s3_error!(InvalidRequest, "iam not init"));
};
let incoming_updated_at = item.updated_at;
let deletion_mark_entities = iam_item_deletion_mark_entities(&item);
match item.r#type.as_str() {
"policy" => apply_iam_policy_item(&iam_sys, &item.name, item.policy).await,
"policy-mapping" => apply_iam_policy_mapping_item(&iam_sys, item.policy_mapping).await,
"group-info" => apply_iam_group_info_item(&iam_sys, item.group_info).await,
let verdict = match item.r#type.as_str() {
"policy" => apply_iam_policy_item(&iam_sys, &item.name, item.policy, incoming_updated_at).await?,
"policy-mapping" => apply_iam_policy_mapping_item(&iam_sys, item.policy_mapping, incoming_updated_at).await?,
"group-info" => apply_iam_group_info_item(&iam_sys, item.group_info, incoming_updated_at).await?,
// MinIO madmin-go sends `SRIAMItemSTSAcc = "sts-account"`. The legacy alias
// `sts-credential` (emitted by older RustFS releases) stays accepted permanently
// so mixed-version RustFS sites keep replicating STS credentials during rolling
// upgrades; it is a compatibility layer, not temporary code.
SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => apply_iam_sts_account_item(&iam_sys, item.sts_credential).await,
"iam-user" => apply_iam_user_item(&iam_sys, item.iam_user, incoming_updated_at).await,
"service-account" => apply_iam_service_account_item(&iam_sys, item.svc_acc_change, incoming_updated_at).await,
_ => Err(s3_error!(
NotImplemented,
"site replication IAM item type `{}` is not supported",
item.r#type
)),
SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => {
apply_iam_sts_account_item(&iam_sys, item.sts_credential).await?;
IamItemVerdict::Apply
}
"iam-user" => apply_iam_user_item(&iam_sys, item.iam_user, incoming_updated_at).await?,
"service-account" => {
apply_iam_service_account_item(&iam_sys, item.svc_acc_change, incoming_updated_at).await?;
IamItemVerdict::Apply
}
_ => {
return Err(s3_error!(
NotImplemented,
"site replication IAM item type `{}` is not supported",
item.r#type
));
}
};
// A committed deletion leaves no record for the gate to judge later items
// against, so its source timestamp is kept as a mark (backlog#2291). The
// mark is part of applying the deletion: failing here makes the sender
// retry the (idempotent) deletion rather than leave a revoke that a stale
// grant could still undo.
if verdict == IamItemVerdict::Apply
&& let Some(deleted_at) = incoming_updated_at.filter(|_| !deletion_mark_entities.is_empty())
{
commit_iam_deletion_marks(deletion_mark_entities, deleted_at).await?;
}
Ok(())
}
/// The deletion mark consulted by the staleness gate when the targeted record
/// is absent: the newest recorded deletion of any of `entities`. An
/// unreadable state falls back to today's behaviour (no mark, the item is
/// applied) — the gate must not turn a state-object outage into rejected
/// IAM replication.
async fn local_iam_deletion_mark(entities: &[String]) -> Option<OffsetDateTime> {
match load_site_replication_state().await {
Ok(state) => iam_deletion_mark(&state, entities),
Err(err) => {
debug!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
result = "iam_deletion_mark_unavailable",
error = ?err,
"site replication state unreadable; applying IAM item without a deletion mark"
);
None
}
}
}
async fn apply_iam_policy_item(iam_sys: &IamSys<ObjectStore>, name: &str, policy: Option<Value>) -> S3Result<()> {
async fn apply_iam_policy_item(
iam_sys: &IamSys<ObjectStore>,
name: &str,
policy: Option<Value>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<IamItemVerdict> {
// Judge the item against the local document's own timestamp so a delayed
// older body (or delete) cannot overwrite a newer edit; once the document
// is deleted, its deletion mark stands in for it (backlog#2291).
let local_updated_at = match iam_sys.get_policy_doc(name).await {
Ok(doc) => Some(doc.update_date.unwrap_or(OffsetDateTime::UNIX_EPOCH)),
Err(err) if rustfs_iam::error::is_err_no_such_policy(&err) => {
local_iam_deletion_mark(&[iam_policy_deletion_mark_entity(name)]).await
}
Err(err) => return Err(ApiError::from(err).into()),
};
if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale {
return Ok(IamItemVerdict::SkipStale);
}
if let Some(policy) = policy {
let policy: Policy =
serde_json::from_value(policy).map_err(|e| s3_error!(InvalidRequest, "invalid policy body: {}", e))?;
@@ -5797,26 +5897,78 @@ async fn apply_iam_policy_item(iam_sys: &IamSys<ObjectStore>, name: &str, policy
Err(err) => return Err(ApiError::from(err).into()),
}
}
Ok(())
Ok(IamItemVerdict::Apply)
}
async fn apply_iam_policy_mapping_item(iam_sys: &IamSys<ObjectStore>, policy_mapping: Option<SRPolicyMapping>) -> S3Result<()> {
async fn apply_iam_policy_mapping_item(
iam_sys: &IamSys<ObjectStore>,
policy_mapping: Option<SRPolicyMapping>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<IamItemVerdict> {
let Some(mapping) = policy_mapping else {
return Err(s3_error!(InvalidRequest, "policyMapping is required"));
};
let user_type = user_type_from_sr_wire(mapping.user_type).ok_or_else(|| s3_error!(InvalidRequest, "invalid userType"))?;
// Judge the item against the stored mapping's timestamp so a delayed older
// attach (or an older detach, `policy == ""`) cannot overwrite a newer one
// (backlog#2291). A detach removes the mapping outright, so once it is
// gone the detach's deletion mark stands in for the record.
let local_updated_at = match iam_sys
.get_mapped_policy_record(&mapping.user_or_group, user_type, mapping.is_group)
.await
{
Some(record) => Some(record.update_at),
None => {
local_iam_deletion_mark(&[iam_policy_mapping_deletion_mark_entity(
&mapping.user_or_group,
mapping.user_type,
mapping.is_group,
)])
.await
}
};
if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale {
return Ok(IamItemVerdict::SkipStale);
}
iam_sys
.policy_db_set(&mapping.user_or_group, user_type, mapping.is_group, &mapping.policy)
.await
.map_err(ApiError::from)?;
Ok(())
Ok(IamItemVerdict::Apply)
}
async fn apply_iam_group_info_item(iam_sys: &IamSys<ObjectStore>, group_info: Option<SRGroupInfo>) -> S3Result<()> {
async fn apply_iam_group_info_item(
iam_sys: &IamSys<ObjectStore>,
group_info: Option<SRGroupInfo>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<IamItemVerdict> {
let Some(group_info) = group_info else {
return Err(s3_error!(InvalidRequest, "groupInfo is required"));
};
let update = group_info.update_req;
// The record is the group itself: its own timestamp moves on every
// membership or status change, so a delayed older add cannot re-add a
// member a newer removal took out, and a delayed older removal (or group
// delete) cannot undo a newer add (backlog#2291). Once the group is gone
// the marks of its deletion and of its members' removals stand in for it,
// so a stale add cannot re-create it or re-add a removed member.
let local_updated_at = match iam_sys.get_group_info(&update.group).await {
Some(group) => Some(group.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)),
None => {
let entities: Vec<String> = std::iter::once(iam_group_deletion_mark_entity(&update.group))
.chain(
update
.members
.iter()
.map(|member| iam_group_member_deletion_mark_entity(&update.group, member)),
)
.collect();
local_iam_deletion_mark(&entities).await
}
};
if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale {
return Ok(IamItemVerdict::SkipStale);
}
if !group_info_requires_upsert(&update) {
// Idempotent removal: a replayed deletion may find the group or a
// member already gone (deleted here earlier, or the user tombstone
@@ -5831,14 +5983,14 @@ async fn apply_iam_group_info_item(iam_sys: &IamSys<ObjectStore>, group_info: Op
}
}
if members.is_empty() && !update.members.is_empty() {
return Ok(());
return Ok(IamItemVerdict::Apply);
}
match iam_sys.remove_users_from_group(&update.group, members).await {
Ok(_) => {}
Err(err) if rustfs_iam::error::is_err_no_such_group(&err) => {}
Err(err) => return Err(ApiError::from(err).into()),
}
return Ok(());
return Ok(IamItemVerdict::Apply);
}
iam_sys
@@ -5849,7 +6001,7 @@ async fn apply_iam_group_info_item(iam_sys: &IamSys<ObjectStore>, group_info: Op
.set_group_status(&update.group, matches!(update.status, GroupStatus::Enabled))
.await
.map_err(ApiError::from)?;
Ok(())
Ok(IamItemVerdict::Apply)
}
async fn apply_iam_sts_account_item(iam_sys: &IamSys<ObjectStore>, sts_credential: Option<SRSTSCredential>) -> S3Result<()> {
@@ -5891,14 +6043,18 @@ async fn apply_iam_user_item(
iam_sys: &IamSys<ObjectStore>,
iam_user: Option<SRIAMUser>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<()> {
) -> S3Result<IamItemVerdict> {
let Some(user) = iam_user else {
return Err(s3_error!(InvalidRequest, "iamUser is required"));
};
if let Some(local) = iam_sys.get_user(&user.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
// Once the identity is deleted, its deletion mark stands in for the
// record so a stale re-create cannot resurrect it (backlog#2291).
let local_updated_at = match iam_sys.get_user(&user.access_key).await {
Some(local) => Some(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)),
None => local_iam_deletion_mark(&[iam_user_deletion_mark_entity(&user.access_key)]).await,
};
if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale {
return Ok(IamItemVerdict::SkipStale);
}
if user.is_delete_req {
iam_sys.delete_user(&user.access_key, true).await.map_err(ApiError::from)?;
@@ -5919,7 +6075,7 @@ async fn apply_iam_user_item(
.map_err(ApiError::from)?;
}
}
Ok(())
Ok(IamItemVerdict::Apply)
}
async fn apply_iam_service_account_item(
@@ -5932,10 +6088,14 @@ async fn apply_iam_service_account_item(
};
let envelope = change.oidc_service_account_envelope;
if let Some(create) = change.create {
let local_updated_at = iam_sys
.get_user(&create.access_key)
.await
.map(|local| local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH));
// Like the user path: with the account already deleted here, the
// recorded deletion mark is the timestamp a stale create/update
// (a snapshot or a delayed delivery) has to beat (backlog#2291).
let local_updated_at = match iam_sys.get_user(&create.access_key).await {
Some(local) => Some(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)),
None if create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT => None,
None => local_iam_deletion_mark(&[format!("svc-acc:{}", create.access_key)]).await,
};
let replicated_policy = if create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT {
if local_updated_at.is_some_and(|local_updated_at| is_stale_update(local_updated_at, incoming_updated_at)) {
return Ok(());
@@ -5979,6 +6139,8 @@ async fn apply_iam_service_account_item(
.map_err(ApiError::from)?;
}
Err(err) if is_err_no_such_service_account(&err) => {
let access_key = create.access_key.clone();
let status = create.status.clone();
iam_sys
.new_service_account(
&create.parent,
@@ -5996,6 +6158,28 @@ async fn apply_iam_service_account_item(
)
.await
.map_err(ApiError::from)?;
// A snapshot (bootstrap / repair / retry resend) carries the
// account's current status; creation always enables, so a
// disabled account must be switched off in a second step or
// the peer keeps accepting credentials the source rejects.
if !status.is_empty() && status != "on" {
iam_sys
.update_service_account(
&access_key,
UpdateServiceAccountOpts {
session_policy: None,
secret_key: None,
name: None,
description: None,
expiration: None,
status: Some(status),
parent_user: None,
allow_site_replicator_account: false,
},
)
.await
.map_err(ApiError::from)?;
}
}
Err(err) => return Err(ApiError::from(err).into()),
}
@@ -6111,6 +6295,7 @@ fn adopt_add_commit_state(state: &mut SiteReplicationState, next_state: SiteRepl
sync_state_initialized,
edit_generation: _,
applied_edit_generations: _,
iam_deletion_marks: _,
} = next_state;
state.name = name;
state.service_account_access_key = service_account_access_key;
@@ -7692,7 +7877,7 @@ impl Operation for SiteReplicationRepairHandler {
let local_peer = current_local_peer(&req, &state);
let body: SiteReplicationRepairRequest = read_site_replication_json(req, "", false).await?;
let info = build_sr_info(&state, &local_peer).await?;
let plan = site_replication_bootstrap_plan(&info)?;
let plan = build_site_replication_bootstrap_plan(&info).await?;
let signing_key = current_token_signing_key().ok_or_else(|| {
S3Error::with_message(S3ErrorCode::InternalError, "token signing key is not initialized".to_string())
})?;
@@ -7885,6 +8070,7 @@ impl Operation for SRRotateServiceAccountHandler {
mod tests {
use super::*;
use crate::site_replication::identity::deployment_id_for_endpoint;
use rustfs_madmin::SRSessionPolicy;
/// A peer the status probe could not reach must render as offline.
///
@@ -12186,6 +12372,299 @@ mod tests {
assert!(!is_stale_update(local, None));
}
/// Minimal model of one replicated IAM record (a policy document body, a
/// user/group mapping, or a group's member set) as the apply paths treat
/// it: `None` is "absent", `Some((content, stamp))` is the local record
/// with the timestamp of the change that last wrote it. Applying an item
/// goes through `judge_iam_item_staleness` exactly like the three apply
/// functions do; a delete (`incoming == None`) on an absent record is the
/// idempotent no-op of backlog#2071.
fn apply_iam_item_to_model(
record: &mut Option<(&'static str, OffsetDateTime)>,
incoming: Option<&'static str>,
incoming_updated_at: Option<OffsetDateTime>,
) -> IamItemVerdict {
let verdict = judge_iam_item_staleness(record.map(|(_, stamp)| stamp), incoming_updated_at);
if verdict == IamItemVerdict::Apply {
*record = incoming.map(|content| (content, incoming_updated_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)));
}
verdict
}
fn at(seconds: i64) -> OffsetDateTime {
OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(seconds)
}
/// backlog#2291: a revoke (narrowed policy body, detached mapping, member
/// removed from the group) followed by the delayed delivery of the older
/// grant must leave the revoke in place.
#[test]
fn test_iam_item_stale_grant_after_revoke_is_not_applied() {
let mut record = None;
assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(10))), IamItemVerdict::Apply);
assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(20))), IamItemVerdict::Apply);
// The older grant is redelivered (retry drain, slow peer) after the revoke.
assert_eq!(
apply_iam_item_to_model(&mut record, Some("grant"), Some(at(10))),
IamItemVerdict::SkipStale,
"a grant older than the local revoke must be acknowledged without being applied"
);
assert_eq!(record, Some(("revoke", at(20))), "the revoke must survive the stale grant");
}
/// backlog#2291: the mirror image — a grant followed by the delayed delivery
/// of an older revoke (older body, older detach, older member removal, or
/// an older delete) must leave the grant in place.
#[test]
fn test_iam_item_stale_revoke_after_grant_is_not_applied() {
let mut record = None;
assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(10))), IamItemVerdict::Apply);
assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(20))), IamItemVerdict::Apply);
assert_eq!(
apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(10))),
IamItemVerdict::SkipStale,
"a revoke older than the local grant must not be applied"
);
assert_eq!(
apply_iam_item_to_model(&mut record, None, Some(at(15))),
IamItemVerdict::SkipStale,
"a delete older than the local record must not remove it"
);
assert_eq!(record, Some(("grant", at(20))));
}
/// backlog#2291: an item at least as new as the local record is applied,
/// including a newer delete; equal timestamps are not stale (same rule as
/// `iam-user`).
#[test]
fn test_iam_item_newer_than_local_record_is_applied() {
let mut record = Some(("grant", at(20)));
assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(20))), IamItemVerdict::Apply);
assert_eq!(record, Some(("revoke", at(20))));
assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(30))), IamItemVerdict::Apply);
assert_eq!(record, Some(("grant", at(30))));
assert_eq!(apply_iam_item_to_model(&mut record, None, Some(at(40))), IamItemVerdict::Apply);
assert_eq!(record, None, "a newer delete removes the record");
}
/// backlog#2291: peers that predate item timestamps keep today's
/// last-writer-wins behaviour — an item without `updatedAt` is applied even
/// over a newer local record.
#[test]
fn test_iam_item_without_source_timestamp_is_applied() {
let mut record = Some(("grant", at(20)));
assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), None), IamItemVerdict::Apply);
assert_eq!(record.map(|(content, _)| content), Some("revoke"));
assert_eq!(judge_iam_item_staleness(Some(at(20)), None), IamItemVerdict::Apply);
assert_eq!(judge_iam_item_staleness(None, None), IamItemVerdict::Apply);
}
/// backlog#2291: nothing is stale relative to an absent record. A create
/// with any timestamp is applied, and a delete falls through to the
/// idempotent no-op paths (backlog#2071) instead of being judged.
#[test]
fn test_iam_item_targeting_absent_record_is_applied() {
assert_eq!(judge_iam_item_staleness(None, Some(at(1))), IamItemVerdict::Apply);
let mut record = None;
assert_eq!(apply_iam_item_to_model(&mut record, None, Some(at(1))), IamItemVerdict::Apply);
assert_eq!(record, None);
assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(1))), IamItemVerdict::Apply);
assert_eq!(record, Some(("grant", at(1))));
// A record that predates timestamps is reported as UNIX_EPOCH by the
// apply paths and therefore never rejects an item.
assert_eq!(
judge_iam_item_staleness(Some(OffsetDateTime::UNIX_EPOCH), Some(at(1))),
IamItemVerdict::Apply
);
}
/// The apply paths once the record is gone: the local timestamp the gate
/// sees is the deletion mark (or `None` when no deletion was recorded),
/// and a committed deletion records its source timestamp as the mark —
/// the same sequence `apply_iam_item` runs.
fn apply_iam_item_to_deleted_record_model(
marks: &mut SiteReplicationState,
entity: &str,
incoming_is_delete: bool,
incoming_updated_at: Option<OffsetDateTime>,
) -> IamItemVerdict {
let entities = vec![entity.to_string()];
let verdict = judge_iam_item_staleness(iam_deletion_mark(marks, &entities), incoming_updated_at);
if verdict == IamItemVerdict::Apply
&& incoming_is_delete
&& let Some(deleted_at) = incoming_updated_at
{
record_iam_deletion_marks(marks, &entities, deleted_at);
}
verdict
}
/// backlog#2291 (real-VM case R6.3a of backlog#2080): a detach deletes the
/// mapping outright, so the older grant that arrives afterwards finds no
/// record — the deletion mark must stand in for it and reject the grant.
/// The same holds for a deleted policy document, user or group.
#[test]
fn test_iam_item_stale_grant_after_record_deletion_is_not_applied() {
let mut marks = SiteReplicationState::default();
let entity = "policy-mapping:alice:0:false";
// The revoke (detach) is applied first: the record is gone, the mark stays.
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))),
IamItemVerdict::Apply
);
assert_eq!(marks.iam_deletion_marks.get(entity), Some(&at(20)));
// The older grant is delivered after the revoke.
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(10))),
IamItemVerdict::SkipStale,
"a grant older than the recorded deletion must not re-create the record"
);
// A replayed copy of the same revoke stays a no-op and keeps the mark.
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))),
IamItemVerdict::Apply
);
assert_eq!(marks.iam_deletion_marks.get(entity), Some(&at(20)));
// An older replayed revoke is stale against the newer one.
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(5))),
IamItemVerdict::SkipStale
);
assert_eq!(
marks.iam_deletion_marks.get(entity),
Some(&at(20)),
"an older deletion never lowers the mark"
);
}
/// backlog#2291: a mark only fences items older than the deletion. A grant
/// newer than (or as new as) the recorded deletion re-creates the record,
/// an unmarked entity and an item without a source timestamp keep today's
/// behaviour.
#[test]
fn test_iam_item_newer_than_deletion_mark_is_applied() {
let mut marks = SiteReplicationState::default();
let entity = "policy:readonly";
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))),
IamItemVerdict::Apply
);
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(20))),
IamItemVerdict::Apply,
"a grant as new as the deletion is not stale"
);
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(30))),
IamItemVerdict::Apply
);
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, entity, false, None),
IamItemVerdict::Apply,
"an item from a peer without timestamps keeps last-writer-wins"
);
assert_eq!(
apply_iam_item_to_deleted_record_model(&mut marks, "policy:other", false, Some(at(1))),
IamItemVerdict::Apply,
"no mark, no record: nothing to be stale against"
);
assert_eq!(
judge_iam_item_staleness(iam_deletion_mark(&marks, &[]), Some(at(1))),
IamItemVerdict::Apply
);
}
/// backlog#2291: a group's removal marks are per member (plus the group
/// itself for a group delete), so with the group gone a stale add is
/// judged against the newest mark among the group and the members it
/// would add.
#[test]
fn test_iam_group_item_after_deletion_is_judged_against_member_marks() {
let mut marks = SiteReplicationState::default();
let bob = iam_group_member_deletion_mark_entity("devs", "bob");
let group = iam_group_deletion_mark_entity("devs");
record_iam_deletion_marks(&mut marks, std::slice::from_ref(&bob), at(20));
record_iam_deletion_marks(&mut marks, std::slice::from_ref(&group), at(30));
// The gate for an add of `bob` to the (deleted) group.
let add_bob = [group.clone(), bob.clone()];
assert_eq!(iam_deletion_mark(&marks, &add_bob), Some(at(30)));
assert_eq!(
judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(25))),
IamItemVerdict::SkipStale
);
assert_eq!(
judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(30))),
IamItemVerdict::Apply
);
// An add of `carol` to a group that was only ever partially emptied
// (no group delete) is judged against carol's own mark only.
marks.iam_deletion_marks.remove(&group);
let add_carol = [group.clone(), iam_group_member_deletion_mark_entity("devs", "carol")];
assert_eq!(iam_deletion_mark(&marks, &add_carol), None);
assert_eq!(
judge_iam_item_staleness(iam_deletion_mark(&marks, &add_carol), Some(at(1))),
IamItemVerdict::Apply
);
let add_bob = [group, bob];
assert_eq!(
judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(10))),
IamItemVerdict::SkipStale
);
}
/// Cheap wiring guard for backlog#2291: every one of the `policy`,
/// `policy-mapping`, `group-info` and `iam-user` apply paths must route
/// through the shared staleness verdict before it writes or deletes
/// anything, and must fall back to the deletion mark when the record is
/// absent; `apply_iam_item` must record the mark of a committed deletion.
/// The ordering rule itself is covered by the `test_iam_item_*` behaviour
/// tests above; this only pins that no path bypasses it again.
#[test]
fn test_iam_policy_mapping_and_group_items_gate_on_incoming_updated_at() {
let source = include_str!("site_replication.rs");
for (start, end) in [
("async fn apply_iam_policy_item(", "async fn apply_iam_policy_mapping_item("),
("async fn apply_iam_policy_mapping_item(", "async fn apply_iam_group_info_item("),
("async fn apply_iam_group_info_item(", "async fn apply_iam_sts_account_item("),
("async fn apply_iam_user_item(", "async fn apply_iam_service_account_item("),
] {
let body = source
.split(start)
.nth(1)
.and_then(|rest| rest.split(end).next())
.expect(start);
assert!(
body.contains("judge_iam_item_staleness(local_updated_at, incoming_updated_at)"),
"{start} must judge the item against the local record before applying it"
);
assert!(
body.contains("local_iam_deletion_mark("),
"{start} must fall back to the deletion mark when the record is absent"
);
}
let dispatch = source
.split("async fn apply_iam_item(")
.nth(1)
.and_then(|rest| rest.split("async fn local_iam_deletion_mark(").next())
.expect("apply_iam_item");
assert!(
dispatch.contains("commit_iam_deletion_marks(deletion_mark_entities, deleted_at)"),
"apply_iam_item must record the mark of a deletion it committed"
);
}
#[test]
fn test_apply_state_edit_req_only_updates_ilm_expiry_flags() {
let mut state = SiteReplicationState::default();
@@ -13717,4 +14196,72 @@ mod tests {
);
}
}
/// backlog#2292: the receiver persists the SOURCE `updated_at` of an
/// applied bucket config and judges the next item's source time against
/// it. Stamping the local apply time instead rejected a source edit that
/// was newer than the applied one but delivered after the local stamp
/// (two quick source edits under delivery delay; a peer clock ahead of
/// ours) and acknowledged it with 200.
#[test]
fn test_bucket_meta_staleness_is_judged_against_the_applied_source_timestamp() {
let apply_wall_clock = OffsetDateTime::now_utc();
let source_edit_t1 = apply_wall_clock - time::Duration::seconds(30);
let source_edit_t2 = source_edit_t1 + time::Duration::seconds(2);
let source_edit_t0 = source_edit_t1 - time::Duration::seconds(2);
assert!(
source_edit_t2 < apply_wall_clock,
"T2 is newer at the source yet older than the local apply clock"
);
// Edit T1 arrives first and is applied the way apply_bucket_meta_item
// persists a replicated config: stamped with its source time.
let mut meta = crate::admin::storage_api::bucket::metadata::BucketMetadata::new("photos");
meta.update_config_at(
BUCKET_POLICY_CONFIG,
br#"{"Version":"2012-10-17","Statement":[]}"#.to_vec(),
source_edit_t1,
)
.expect("apply edit T1");
let local_updated_at = bucket_meta_local_updated_at(&meta, BUCKET_POLICY_CONFIG);
assert_eq!(
local_updated_at, source_edit_t1,
"the stored stamp is the source time, not the apply clock"
);
// Edit T2 is newer at the source but delivered late: it must apply.
assert!(
!is_stale_update(local_updated_at, Some(source_edit_t2)),
"edit T2 ({source_edit_t2}) is newer than applied edit T1 ({source_edit_t1}) but is rejected against local stamp {local_updated_at}"
);
// Edit T0 predates the applied edit: it stays rejected.
assert!(
is_stale_update(local_updated_at, Some(source_edit_t0)),
"edit T0 ({source_edit_t0}) is older than applied edit T1 ({source_edit_t1}) and must be rejected"
);
// An item without a source time is never judged stale (unchanged).
assert!(!is_stale_update(local_updated_at, None));
}
/// backlog#2292: the replicated-config write in `apply_bucket_meta_item`
/// must go through the source-stamped entries; a plain
/// `update_if_incarnation` there would reintroduce local stamping.
#[test]
fn test_apply_bucket_meta_item_writes_through_the_source_stamped_entries() {
let source = include_str!("site_replication.rs");
let apply = source
.split("async fn apply_bucket_meta_item")
.nth(1)
.and_then(|rest| rest.split("fn group_info_requires_upsert").next())
.expect("apply_bucket_meta_item source");
assert!(
apply.contains("update_quota_if_incarnation_at("),
"durable quota must carry the source stamp"
);
assert!(apply.contains("update_if_incarnation_at("), "bucket configs must carry the source stamp");
assert!(
!apply.contains("metadata_sys::update_if_incarnation(&item.bucket"),
"no replicated config write may bypass the source stamp"
);
}
}
+38
View File
@@ -353,6 +353,25 @@ pub(crate) mod metadata_sys {
super::ecstore_bucket::metadata_sys::update_if_incarnation(bucket, config_file, data, expected_incarnation_id).await
}
/// [`update_if_incarnation`] stamping the config with a replicated edit's
/// source `updated_at` instead of the local clock (backlog#2292).
pub(crate) async fn update_if_incarnation_at(
bucket: &str,
config_file: &str,
data: Vec<u8>,
expected_incarnation_id: uuid::Uuid,
updated_at: OffsetDateTime,
) -> Result<OffsetDateTime> {
super::ecstore_bucket::metadata_sys::update_if_incarnation_at(
bucket,
config_file,
data,
expected_incarnation_id,
updated_at,
)
.await
}
pub(crate) async fn update_quota_if_incarnation(
bucket: &str,
data: Vec<u8>,
@@ -362,6 +381,25 @@ pub(crate) mod metadata_sys {
super::ecstore_bucket::metadata_sys::update_quota_if_incarnation(bucket, data, expected_incarnation_id, proof).await
}
/// [`update_quota_if_incarnation`] stamping the quota with a replicated
/// edit's source `updated_at` instead of the local clock (backlog#2292).
pub(crate) async fn update_quota_if_incarnation_at(
bucket: &str,
data: Vec<u8>,
expected_incarnation_id: uuid::Uuid,
proof: &super::ecstore_notification::CrossPoolFenceFleetProofToken,
updated_at: OffsetDateTime,
) -> Result<OffsetDateTime> {
super::ecstore_bucket::metadata_sys::update_quota_if_incarnation_at(
bucket,
data,
expected_incarnation_id,
proof,
updated_at,
)
.await
}
pub(crate) async fn capture_bucket_metadata_incarnation(bucket: &str) -> Result<uuid::Uuid> {
super::ecstore_bucket::metadata_sys::capture_bucket_metadata_incarnation(bucket).await
}
+2
View File
@@ -726,6 +726,8 @@ pub(crate) mod bucket {
delete_marker_version_id: None,
delete_marker: false,
delete_marker_mtime: None,
target_delete_marker_version_ids: Default::default(),
target_delete_marker_version_ids_corrupt: false,
target_arns,
force_delete_id: Some(operation_id),
force_delete_generation: Some(i64::try_from(generation.unix_timestamp_nanos()).unwrap_or(i64::MAX)),
+229 -18
View File
@@ -302,7 +302,164 @@ pub(crate) fn site_replication_state_replicates_ilm_expiry(state: &SiteReplicati
state.peers.values().any(|peer| peer.replicate_ilm_expiry)
}
pub(crate) fn site_replication_bootstrap_plan(info: &SRInfo) -> S3Result<SiteReplicationBootstrapPlan> {
/// Secret-bearing half of the IAM snapshot. `SRInfo` is served to admin
/// callers (`site-replication/info`, status, add preflight) and must stay
/// secret-free, so the bootstrap plan receives credentials through this
/// separate value, built only on the paths that deliver to peers (site add
/// bootstrap, repair, retry snapshot resend). Never persisted, never served.
#[derive(Debug, Clone, Default)]
pub(crate) struct SiteReplicationIamCredentials {
/// Built-in users (access key -> credential); temp and service accounts
/// are excluded, external/IdP users never appear here.
pub(crate) users: BTreeMap<String, SiteReplicationUserCredential>,
/// Every service account except the site replicator's own, already
/// shaped as the `service-account` create item the live hook emits.
pub(crate) service_accounts: Vec<SiteReplicationServiceAccountSnapshot>,
}
#[derive(Debug, Clone)]
pub(crate) struct SiteReplicationUserCredential {
pub(crate) secret_key: String,
pub(crate) status: AccountStatus,
/// The user record's own update time (the axis the receiver's staleness
/// check compares against), unlike `UserInfo::updated_at` which
/// `list_users` overwrites with the policy mapping's time.
pub(crate) updated_at: Option<OffsetDateTime>,
}
#[derive(Debug, Clone)]
pub(crate) struct SiteReplicationServiceAccountSnapshot {
pub(crate) create: SRSvcAccCreate,
pub(crate) envelope: Option<SRSvcAccReplicationEnvelope>,
pub(crate) updated_at: Option<OffsetDateTime>,
}
pub(crate) const SERVICE_ACCOUNT_ENVELOPE_VERSION: u64 = 2;
pub(crate) fn encode_service_account_replication_policy(
claims: &HashMap<String, Value>,
session_policy: Option<&str>,
) -> S3Result<(SRSessionPolicy, Option<SRSvcAccReplicationEnvelope>)> {
if !claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM) {
return session_policy
.map(SRSessionPolicy::from_json)
.transpose()
.map(|policy| policy.unwrap_or_default())
.map(|policy| (policy, None))
.map_err(|err| s3_error!(InvalidArgument, "marshal policy failed: {:?}", err));
}
let policy = match session_policy {
Some(policy) => serde_json::from_str::<Policy>(policy)
.map_err(|err| s3_error!(InvalidArgument, "invalid service account replication policy: {:?}", err))?,
None => Policy::default(),
};
if policy.statements.is_empty() && (!policy.id.is_empty() || !policy.version.is_empty())
|| policy.version.is_empty() && !policy.statements.is_empty()
{
return Err(s3_error!(InvalidArgument, "service account replication policy is not normalized"));
}
let policy = serde_json::to_string(&policy)
.map_err(|err| s3_error!(InternalError, "marshal service account replication policy failed: {:?}", err))?;
let policy = SRSessionPolicy::from_json(&policy)
.map_err(|err| s3_error!(InternalError, "marshal service account replication policy failed: {:?}", err))?;
Ok((
policy,
Some(SRSvcAccReplicationEnvelope {
version: SERVICE_ACCOUNT_ENVELOPE_VERSION,
}),
))
}
/// Read the credentials the IAM snapshot needs straight from the IAM store:
/// `list_users` deliberately strips secret keys and skips service accounts,
/// which is right for an admin listing and wrong for a peer snapshot (the
/// plan builder used to drop every user for lack of a secret, so a status
/// change or secret rotation committed while a peer was unreachable never
/// reached it — backlog#2289).
pub(crate) async fn build_sr_iam_credentials() -> S3Result<SiteReplicationIamCredentials> {
let mut credentials = SiteReplicationIamCredentials::default();
let Some(iam_sys) = current_iam_handle() else {
return Ok(credentials);
};
let mut users = HashMap::new();
iam_sys.load_users(UserType::Reg, &mut users).await.map_err(ApiError::from)?;
for (access_key, identity) in users {
if identity.credentials.is_temp() || identity.credentials.is_service_account() {
continue;
}
credentials.users.insert(
access_key,
SiteReplicationUserCredential {
secret_key: identity.credentials.secret_key,
status: if identity.credentials.status == "off" {
AccountStatus::Disabled
} else {
AccountStatus::Enabled
},
updated_at: identity.update_at,
},
);
}
let mut service_accounts = HashMap::new();
iam_sys
.load_users(UserType::Svc, &mut service_accounts)
.await
.map_err(ApiError::from)?;
let mut service_accounts: Vec<_> = service_accounts.into_iter().collect();
service_accounts.sort_by(|(a, _), (b, _)| a.cmp(b));
for (access_key, identity) in service_accounts {
// The replicator account is installed by join / rotate, never by a snapshot.
if access_key == SITE_REPLICATOR_SERVICE_ACCOUNT || !identity.credentials.is_service_account() {
continue;
}
let claims = iam_sys.get_claims_for_svc_acc(&access_key).await.map_err(ApiError::from)?;
let (account, session_policy) = iam_sys.get_service_account(&access_key).await.map_err(ApiError::from)?;
let session_policy = session_policy
.map(|policy| serde_json::to_string(&policy))
.transpose()
.map_err(|err| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("marshal service account session policy failed: {err:?}"),
)
})?;
let (session_policy, envelope) = encode_service_account_replication_policy(&claims, session_policy.as_deref())?;
credentials.service_accounts.push(SiteReplicationServiceAccountSnapshot {
create: SRSvcAccCreate {
parent: identity.credentials.parent_user,
access_key,
secret_key: identity.credentials.secret_key,
groups: identity.credentials.groups.unwrap_or_default(),
claims,
session_policy,
status: identity.credentials.status,
name: account.name.unwrap_or_default(),
description: account.description.unwrap_or_default(),
expiration: account.expiration,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
},
envelope,
updated_at: identity.update_at,
});
}
Ok(credentials)
}
/// The bootstrap plan for peer delivery: `info` (secret-free) plus the IAM
/// credentials read at this moment.
pub(crate) async fn build_site_replication_bootstrap_plan(info: &SRInfo) -> S3Result<SiteReplicationBootstrapPlan> {
let credentials = build_sr_iam_credentials().await?;
site_replication_bootstrap_plan(info, &credentials)
}
pub(crate) fn site_replication_bootstrap_plan(
info: &SRInfo,
credentials: &SiteReplicationIamCredentials,
) -> S3Result<SiteReplicationBootstrapPlan> {
let mut plan = SiteReplicationBootstrapPlan::default();
let replicate_ilm_expiry = site_replication_info_replicates_ilm_expiry(info);
@@ -318,24 +475,57 @@ pub(crate) fn site_replication_bootstrap_plan(info: &SRInfo) -> S3Result<SiteRep
}
for (access_key, user) in &info.user_info_map {
if let Some(secret_key) = &user.secret_key {
plan.iam_items.push(SRIAMItem {
r#type: "iam-user".to_string(),
iam_user: Some(rustfs_madmin::SRIAMUser {
access_key: access_key.clone(),
is_delete_req: false,
user_req: Some(AddOrUpdateUserReq {
secret_key: secret_key.clone(),
policy: user.policy_name.clone(),
status: user.status.clone(),
}),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
// Credentials come from the store snapshot; an inline `secret_key` on
// the SRInfo entry (older callers, tests) is accepted as a fallback.
// Users with neither (external / IdP identities) have nothing a peer
// could install and are skipped.
let credential = credentials.users.get(access_key);
let Some(secret_key) = credential
.map(|credential| credential.secret_key.clone())
.or_else(|| user.secret_key.clone())
.filter(|secret_key| !secret_key.is_empty())
else {
continue;
};
let status = credential
.map(|credential| credential.status.clone())
.unwrap_or_else(|| user.status.clone());
let updated_at = credential.and_then(|credential| credential.updated_at).or(user.updated_at);
plan.iam_items.push(SRIAMItem {
r#type: "iam-user".to_string(),
iam_user: Some(rustfs_madmin::SRIAMUser {
access_key: access_key.clone(),
is_delete_req: false,
user_req: Some(AddOrUpdateUserReq {
secret_key,
policy: user.policy_name.clone(),
status,
}),
updated_at: user.updated_at,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
}),
updated_at,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
});
}
// Service accounts follow their parents: the receiver creates a missing
// account under `parent` and updates an existing one (secret, status,
// session policy), so a rotation or disable committed during an outage
// converges through the same snapshot as users do.
for account in &credentials.service_accounts {
plan.iam_items.push(SRIAMItem {
r#type: "service-account".to_string(),
svc_acc_change: Some(SRSvcAccChange {
create: Some(account.create.clone()),
oidc_service_account_envelope: account.envelope.clone(),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
});
}
}),
updated_at: account.updated_at,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
});
}
for (name, desc) in &info.group_desc_map {
@@ -518,7 +708,12 @@ pub(crate) async fn broadcast_site_replication_make_bucket(
} else {
path
};
broadcast_site_replication_json_using_runtime(runtime, &path, &serde_json::json!({})).await?;
// Both steps run to completion on their own: the broadcast attempts every
// peer and reports the first failure (backlog#2293), so stopping here on
// that error would skip `configure-replication` for the peers whose
// `make` just succeeded — and nothing records a retry for that gap. The
// failed peer's retry events cover both steps independently.
let make_result = broadcast_site_replication_json_using_runtime(runtime, &path, &serde_json::json!({})).await;
let configure_path = bootstrap_bucket_op_path(bucket, "configure-replication");
let configure_path = if let Some(token) = bootstrap_token {
@@ -526,7 +721,8 @@ pub(crate) async fn broadcast_site_replication_make_bucket(
} else {
configure_path
};
broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await
let configure_result = broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await;
make_result.and(configure_result)
}
const SITE_REPLICATION_DELETE_INTENT_PENDING: &str =
@@ -832,6 +1028,21 @@ pub async fn site_replication_iam_change_hook(item: SRIAMItem) -> S3Result<()> {
let Some(runtime) = runtime_site_replication_targets().await? else {
return Ok(());
};
// A local revoke must out-rank a stale grant a peer delivers later, so its
// mark is committed before the broadcast (backlog#2291). The broadcast
// still goes out when the mark cannot be persisted: the peers' own records
// remain the primary gate, the mark only covers the deleted case.
if let Err(err) = record_iam_deletion_marks_for_item(&item).await {
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
item_type = %item.r#type,
result = "iam_deletion_mark_not_recorded",
error = ?err,
"failed to record local IAM deletion mark before broadcast"
);
}
let mut first_error: Option<S3Error> = None;
for peer in runtime.state.peers.values() {
if peer.deployment_id == runtime.local_peer.deployment_id
+26 -3
View File
@@ -79,13 +79,16 @@ use http::header::{CONTENT_TYPE, HOST};
use http::{HeaderMap, HeaderValue, Uri};
use hyper::{Method, StatusCode};
use rustfs_config::{DEFAULT_CONSOLE_ADDRESS, DEFAULT_RUSTFS_TLS_PATH, ENV_RUSTFS_CONSOLE_ADDRESS, ENV_RUSTFS_TLS_PATH};
use rustfs_iam::federation::OIDC_VIRTUAL_PARENT_CLAIM;
use rustfs_iam::store::{MappedPolicy, UserType, sr_wire_user_type};
use rustfs_iam::sys::SITE_REPLICATOR_SERVICE_ACCOUNT;
use rustfs_madmin::{
AddOrUpdateUserReq, GroupAddRemove, GroupStatus, PeerInfo, PeerSite, ReplicateEditStatus, SITE_REPL_API_VERSION,
SRBucketInfo, SRBucketMeta, SRGroupInfo, SRIAMItem, SRIAMPolicy, SRInfo, SRPolicyMapping, SRRemoveReq, SRResyncOpStatus,
SRRetryStats, SRStateInfo, SyncStatus,
AccountStatus, AddOrUpdateUserReq, GroupAddRemove, GroupStatus, PeerInfo, PeerSite, ReplicateEditStatus,
SITE_REPL_API_VERSION, SRBucketInfo, SRBucketMeta, SRGroupInfo, SRIAMItem, SRIAMPolicy, SRInfo, SRPolicyMapping, SRRemoveReq,
SRResyncOpStatus, SRRetryStats, SRSessionPolicy, SRStateInfo, SRSvcAccChange, SRSvcAccCreate, SRSvcAccDelete,
SRSvcAccReplicationEnvelope, SyncStatus,
};
use rustfs_policy::policy::Policy;
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
use rustfs_signer::sign_v4;
use rustfs_tls_runtime::{GlobalPublishedOutboundTlsState, TlsGeneration};
@@ -107,6 +110,26 @@ use tracing::{info, warn};
use url::{Url, form_urlencoded};
use uuid::Uuid;
/// Serialize `value` with every JSON object's keys sorted, for hashing and
/// equality checks. `HashMap` fields (service-account claims) iterate in a
/// per-instance random order and `serde_json` is built with `preserve_order`,
/// so two identical plans would otherwise hash differently: the repair
/// preflight token went stale between dry-run and execute, and a retry
/// snapshot resend never looked "stable" (backlog#2289 follow-up).
pub(crate) fn canonical_json_vec<T: Serialize>(value: &T) -> serde_json::Result<Vec<u8>> {
fn sort_keys(value: Value) -> Value {
match value {
Value::Object(map) => {
let sorted: BTreeMap<String, Value> = map.into_iter().map(|(key, value)| (key, sort_keys(value))).collect();
Value::Object(sorted.into_iter().collect())
}
Value::Array(items) => Value::Array(items.into_iter().map(sort_keys).collect()),
other => other,
}
}
serde_json::to_vec(&sort_keys(serde_json::to_value(value)?))
}
pub(crate) const LOG_COMPONENT_ADMIN: &str = "admin";
pub(crate) const LOG_SUBSYSTEM_SITE_REPLICATION: &str = "site_replication";
+3 -3
View File
@@ -234,9 +234,9 @@ impl SiteReplicationRepairTask<'_> {
pub(crate) fn id(&self) -> S3Result<String> {
let payload = match self {
Self::Iam(item) => serde_json::to_vec(item),
Self::Iam(item) => canonical_json_vec(item),
Self::BucketMake(_) | Self::Replication(_) => serde_json::to_vec(&serde_json::json!({})),
Self::BucketMetadata(item) => serde_json::to_vec(item),
Self::BucketMetadata(item) => canonical_json_vec(item),
}
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize repair task failed: {err}")))?;
let mut digest = Sha256::new();
@@ -726,7 +726,7 @@ pub(crate) async fn execute_site_replication_repair_locked(
return Err(s3_error!(InvalidRequest, "site replication is not configured"));
}
let info = build_sr_info(&state, &request.local_peer).await?;
let plan = site_replication_bootstrap_plan(&info)?;
let plan = build_site_replication_bootstrap_plan(&info).await?;
let plan_token = site_replication_repair_plan_token(&state, &plan)?;
let preflight_token = site_replication_repair_preflight_token(&state, &plan, request.signing_key.as_bytes())?;
let sites = site_replication_repair_sites(&state, &request.local_peer, &plan, request.signing_key.as_bytes())?;
+107 -7
View File
@@ -397,12 +397,12 @@ pub(crate) fn iam_deletion_replay_matches(record: &SiteReplicationIamDeletionRep
/// newer revision of one another.
pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option<String> {
match item.r#type.as_str() {
"policy" if item.policy.is_none() => Some(format!("policy:{}", item.name)),
"policy" if item.policy.is_none() => Some(iam_policy_deletion_mark_entity(&item.name)),
"iam-user" => item
.iam_user
.as_ref()
.filter(|user| user.is_delete_req)
.map(|user| format!("iam-user:{}", user.access_key)),
.map(|user| iam_user_deletion_mark_entity(&user.access_key)),
"group-info" => item
.group_info
.as_ref()
@@ -416,7 +416,7 @@ pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option<String> {
.policy_mapping
.as_ref()
.filter(|mapping| mapping.policy.is_empty())
.map(|mapping| format!("policy-mapping:{}:{}:{}", mapping.user_or_group, mapping.user_type, mapping.is_group)),
.map(|mapping| iam_policy_mapping_deletion_mark_entity(&mapping.user_or_group, mapping.user_type, mapping.is_group)),
"service-account" => item
.svc_acc_change
.as_ref()
@@ -426,6 +426,82 @@ pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option<String> {
}
}
/// The entities whose deletion a deletion-shaped IAM item commits, keyed the
/// way the receive-side staleness gate looks them up once the local record is
/// gone (backlog#2291); empty for creates and updates. Group member removal
/// yields one entity per removed member so a stale re-add of that member can
/// be judged, and a group delete (no members) yields the group itself.
pub(crate) fn iam_item_deletion_mark_entities(item: &SRIAMItem) -> Vec<String> {
if item.r#type == "group-info" {
let Some(update) = item
.group_info
.as_ref()
.map(|group| &group.update_req)
.filter(|update| update.is_remove)
else {
return Vec::new();
};
if update.members.is_empty() {
return vec![iam_group_deletion_mark_entity(&update.group)];
}
return update
.members
.iter()
.map(|member| iam_group_member_deletion_mark_entity(&update.group, member))
.collect();
}
iam_item_deletion_entity(item).into_iter().collect()
}
pub(crate) fn iam_policy_deletion_mark_entity(name: &str) -> String {
format!("policy:{name}")
}
pub(crate) fn iam_user_deletion_mark_entity(access_key: &str) -> String {
format!("iam-user:{access_key}")
}
/// `user_type` is the SR wire integer, as carried by the item on both sides.
pub(crate) fn iam_policy_mapping_deletion_mark_entity(user_or_group: &str, user_type: i64, is_group: bool) -> String {
format!("policy-mapping:{user_or_group}:{user_type}:{is_group}")
}
pub(crate) fn iam_group_deletion_mark_entity(group: &str) -> String {
format!("group:{group}")
}
pub(crate) fn iam_group_member_deletion_mark_entity(group: &str, member: &str) -> String {
format!("group-member:{group}:{member}")
}
/// Persist the deletion marks of `item` (its source `updated_at` per entity
/// of [`iam_item_deletion_mark_entities`]) through the state transaction.
/// No-op for creates/updates and for items without a source timestamp
/// (older peers): a mark without a source clock could not be ordered against
/// later items. Called before a local deletion is broadcast and after a
/// replicated deletion is applied, so both sides out-rank a stale grant that
/// arrives later.
pub(crate) async fn record_iam_deletion_marks_for_item(item: &SRIAMItem) -> S3Result<()> {
let entities = iam_item_deletion_mark_entities(item);
let Some(deleted_at) = item.updated_at.filter(|_| !entities.is_empty()) else {
return Ok(());
};
commit_iam_deletion_marks(entities, deleted_at).await
}
/// [`record_iam_deletion_marks`] under the state transaction; the write is
/// skipped when no mark moves.
pub(crate) async fn commit_iam_deletion_marks(entities: Vec<String>, deleted_at: OffsetDateTime) -> S3Result<()> {
update_site_replication_state_when_changed(move |state| {
Ok(if record_iam_deletion_marks(state, &entities, deleted_at) {
StateCommit::Changed(())
} else {
StateCommit::Unchanged(())
})
})
.await
}
/// Failure bookkeeping for one IAM item delivery: upsert the collapsed retry
/// event and, when the item is a deletion, record its body for replay. Both
/// live in the same state so the caller commits them in one transaction — a
@@ -791,8 +867,8 @@ impl RetrySnapshot {
pub(crate) fn fingerprint(&self) -> S3Result<Vec<Vec<u8>>> {
let mut payloads = match self {
Self::Iam(items) => items.iter().map(serde_json::to_vec).collect::<Result<Vec<_>, _>>(),
Self::BucketMetadata(items) => items.iter().map(serde_json::to_vec).collect::<Result<Vec<_>, _>>(),
Self::Iam(items) => items.iter().map(canonical_json_vec).collect::<Result<Vec<_>, _>>(),
Self::BucketMetadata(items) => items.iter().map(canonical_json_vec).collect::<Result<Vec<_>, _>>(),
}
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize retry snapshot failed: {err}")))?;
payloads.sort_unstable();
@@ -954,6 +1030,7 @@ pub(crate) enum IamSnapshotKey {
User(String),
Group(String),
PolicyMapping { target: String, user_type: i64, is_group: bool },
ServiceAccount(String),
}
pub(crate) fn iam_snapshot_key(item: &SRIAMItem) -> Option<IamSnapshotKey> {
@@ -972,6 +1049,11 @@ pub(crate) fn iam_snapshot_key(item: &SRIAMItem) -> Option<IamSnapshotKey> {
user_type: mapping.user_type,
is_group: mapping.is_group,
}),
"service-account" => item
.svc_acc_change
.as_ref()
.and_then(|change| change.create.as_ref())
.map(|create| IamSnapshotKey::ServiceAccount(create.access_key.clone())),
_ => None,
}
}
@@ -1006,6 +1088,24 @@ pub(crate) fn iam_snapshot_tombstones(item: &SRIAMItem, observed_at: OffsetDateT
mapping.policy.clear();
}
}
"service-account" => {
let Some(access_key) = item
.svc_acc_change
.as_ref()
.and_then(|change| change.create.as_ref())
.map(|create| create.access_key.clone())
else {
return Vec::new();
};
tombstone.svc_acc_change = Some(SRSvcAccChange {
delete: Some(SRSvcAccDelete {
access_key,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
}),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
});
}
_ => return Vec::new(),
}
vec![tombstone]
@@ -1701,7 +1801,7 @@ pub(crate) async fn drain_site_replication_retry_queue_locked(
// tick and only when a snapshot resend is actually due.
let plan = if needs_plan {
let info = build_sr_info(&runtime.state, &runtime.local_peer).await?;
Some(site_replication_bootstrap_plan(&info)?)
Some(build_site_replication_bootstrap_plan(&info).await?)
} else {
None
};
@@ -1841,7 +1941,7 @@ pub(crate) async fn drain_one_site_replication_retry_event(
}
}
let fresh_info = build_sr_info(&runtime.state, &runtime.local_peer).await?;
let fresh_plan = site_replication_bootstrap_plan(&fresh_info)?;
let fresh_plan = build_site_replication_bootstrap_plan(&fresh_info).await?;
let fresh_snapshot = RetrySnapshot::from_plan(&action, &fresh_plan).expect("snapshot action has a snapshot");
if fresh_snapshot.fingerprint()? == current_fingerprint {
if is_iam {
+80
View File
@@ -64,6 +64,86 @@ pub(crate) struct SiteReplicationState {
/// newer edit that already landed.
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub(crate) applied_edit_generations: BTreeMap<String, u64>,
/// Source timestamp of the newest IAM deletion committed on this site,
/// keyed by the deleted entity (`iam_item_deletion_mark_entities`). A
/// deletion leaves no local record to judge a later item against, so this
/// is what lets the receive-side staleness gate reject a grant that is
/// older than the revoke it would otherwise undo (backlog#2291). Bounded
/// by [`SITE_REPLICATION_IAM_DELETION_MARK_LIMIT`]; the oldest mark is
/// evicted first.
#[serde(default, with = "rfc3339_map", skip_serializing_if = "BTreeMap::is_empty")]
pub(crate) iam_deletion_marks: BTreeMap<String, OffsetDateTime>,
}
/// Upper bound on [`SiteReplicationState::iam_deletion_marks`].
pub(crate) const SITE_REPLICATION_IAM_DELETION_MARK_LIMIT: usize = 1024;
/// Record that deletions of `entities` with source timestamp `deleted_at`
/// were committed here. Newest wins per entity: an older deletion never
/// lowers a mark. Returns whether the state changed.
pub(crate) fn record_iam_deletion_marks(
state: &mut SiteReplicationState,
entities: &[String],
deleted_at: OffsetDateTime,
) -> bool {
let mut changed = false;
for entity in entities {
if state
.iam_deletion_marks
.get(entity)
.is_some_and(|existing| *existing >= deleted_at)
{
continue;
}
state.iam_deletion_marks.insert(entity.clone(), deleted_at);
changed = true;
}
while state.iam_deletion_marks.len() > SITE_REPLICATION_IAM_DELETION_MARK_LIMIT {
let Some(oldest) = state
.iam_deletion_marks
.iter()
.min_by_key(|(_, deleted_at)| **deleted_at)
.map(|(entity, _)| entity.clone())
else {
break;
};
state.iam_deletion_marks.remove(&oldest);
}
changed
}
/// Newest deletion mark among `entities`, or `None` when no deletion of any
/// of them was recorded here. The receive-side staleness gate feeds this in
/// as the local timestamp when the targeted record is absent.
pub(crate) fn iam_deletion_mark(state: &SiteReplicationState, entities: &[String]) -> Option<OffsetDateTime> {
entities
.iter()
.filter_map(|entity| state.iam_deletion_marks.get(entity).copied())
.max()
}
/// RFC 3339 map values, matching the other timestamps in the state object
/// (`time::serde::rfc3339` only applies to a single field).
mod rfc3339_map {
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::collections::BTreeMap;
use time::OffsetDateTime;
#[derive(Serialize, Deserialize)]
#[serde(transparent)]
struct Stamp(#[serde(with = "time::serde::rfc3339")] OffsetDateTime);
pub(super) fn serialize<S: Serializer>(map: &BTreeMap<String, OffsetDateTime>, serializer: S) -> Result<S::Ok, S::Error> {
serializer.collect_map(map.iter().map(|(entity, deleted_at)| (entity, Stamp(*deleted_at))))
}
pub(super) fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result<BTreeMap<String, OffsetDateTime>, D::Error> {
let map = BTreeMap::<String, Stamp>::deserialize(deserializer)?;
Ok(map
.into_iter()
.map(|(entity, Stamp(deleted_at))| (entity, deleted_at))
.collect())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
+465 -5
View File
@@ -554,6 +554,121 @@ fn test_iam_item_deletion_entity_shapes() {
assert!(iam_item_deletion_entity(&policy_set).is_none());
}
/// Deletion marks (backlog#2291) key on the same entities as the replay
/// records, except that a group member removal is marked per member (so a
/// stale re-add of one member can be judged) and a group delete marks the
/// group itself. Creates and updates leave no mark.
#[test]
fn test_iam_item_deletion_mark_entities_shapes() {
assert_eq!(
iam_item_deletion_mark_entities(&user_delete_item("alice")),
vec!["iam-user:alice".to_string()]
);
assert_eq!(
iam_item_deletion_mark_entities(&policy_delete_item("readonly")),
vec!["policy:readonly".to_string()]
);
let mut group_remove = SRIAMItem {
r#type: "group-info".to_string(),
group_info: Some(SRGroupInfo {
update_req: GroupAddRemove {
group: "devs".to_string(),
members: vec!["bob".to_string(), "alice".to_string()],
status: GroupStatus::Enabled,
is_remove: true,
},
api_version: None,
}),
..Default::default()
};
assert_eq!(
iam_item_deletion_mark_entities(&group_remove),
vec!["group-member:devs:bob".to_string(), "group-member:devs:alice".to_string()]
);
group_remove
.group_info
.as_mut()
.expect("group info")
.update_req
.members
.clear();
assert_eq!(
iam_item_deletion_mark_entities(&group_remove),
vec!["group:devs".to_string()],
"a removal without members deletes the group"
);
group_remove.group_info.as_mut().expect("group info").update_req.is_remove = false;
assert!(iam_item_deletion_mark_entities(&group_remove).is_empty());
let mapping_clear = SRIAMItem {
r#type: "policy-mapping".to_string(),
policy_mapping: Some(SRPolicyMapping {
user_or_group: "alice".to_string(),
user_type: 0,
is_group: false,
policy: String::new(),
..Default::default()
}),
..Default::default()
};
assert_eq!(
iam_item_deletion_mark_entities(&mapping_clear),
vec!["policy-mapping:alice:0:false".to_string()]
);
let mut user_create = user_delete_item("alice");
user_create.iam_user.as_mut().expect("iam user").is_delete_req = false;
assert!(iam_item_deletion_mark_entities(&user_create).is_empty());
}
/// Newest wins per entity, the map stays bounded by evicting the oldest
/// mark, and the timestamps survive the state object as RFC 3339.
#[test]
fn test_record_iam_deletion_marks_newest_wins_and_stays_bounded() {
let at = |seconds: i64| OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(seconds);
let mut state = SiteReplicationState::default();
let alice = vec!["iam-user:alice".to_string()];
assert!(record_iam_deletion_marks(&mut state, &alice, at(20)));
assert!(
!record_iam_deletion_marks(&mut state, &alice, at(10)),
"an older deletion does not move the mark"
);
assert!(
!record_iam_deletion_marks(&mut state, &alice, at(20)),
"a replayed deletion is not a change"
);
assert_eq!(iam_deletion_mark(&state, &alice), Some(at(20)));
assert!(record_iam_deletion_marks(&mut state, &alice, at(30)));
assert_eq!(iam_deletion_mark(&state, &alice), Some(at(30)));
assert_eq!(iam_deletion_mark(&state, &["iam-user:bob".to_string()]), None);
assert!(!record_iam_deletion_marks(&mut state, &[], at(40)));
// Fill past the bound with marks older than alice's; the oldest go first.
let members: Vec<String> = (0..SITE_REPLICATION_IAM_DELETION_MARK_LIMIT)
.map(|index| format!("group-member:devs:user-{index:04}"))
.collect();
for (index, member) in members.iter().enumerate() {
record_iam_deletion_marks(&mut state, std::slice::from_ref(member), at(index as i64 - 2000));
}
assert_eq!(state.iam_deletion_marks.len(), SITE_REPLICATION_IAM_DELETION_MARK_LIMIT);
assert_eq!(iam_deletion_mark(&state, &alice), Some(at(30)), "the newest mark survives eviction");
assert_eq!(iam_deletion_mark(&state, &members[..1]), None, "the oldest mark is evicted first");
assert_eq!(iam_deletion_mark(&state, &members[1..2]), Some(at(-1999)));
let json = serde_json::to_value(&state).expect("serialize state");
assert_eq!(json["iam_deletion_marks"]["iam-user:alice"], serde_json::json!("1970-01-01T00:00:30Z"));
let reloaded = parse_site_replication_state(&serde_json::to_vec(&state).expect("serialize state")).expect("parse state");
assert_eq!(reloaded.iam_deletion_marks, state.iam_deletion_marks);
assert!(
parse_site_replication_state(br#"{"name":"a","service_account_access_key":"","service_account_parent":"","peers":{},"updated_at":null,"resync_status":{}}"#)
.expect("state without marks")
.iam_deletion_marks
.is_empty()
);
}
/// 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
@@ -1679,7 +1794,8 @@ fn test_site_replication_bootstrap_plan_includes_replayable_snapshot_items() {
},
);
let plan = site_replication_bootstrap_plan(&info).expect("bootstrap plan should build");
let plan =
site_replication_bootstrap_plan(&info, &SiteReplicationIamCredentials::default()).expect("bootstrap plan should build");
assert_eq!(plan.iam_items.iter().map(|item| item.r#type.as_str()).collect::<Vec<_>>(), {
vec!["policy", "iam-user", "group-info", "policy-mapping"]
@@ -1717,7 +1833,8 @@ fn test_site_replication_bootstrap_plan_skips_lifecycle_by_default() {
},
);
let plan = site_replication_bootstrap_plan(&info).expect("bootstrap plan should build");
let plan =
site_replication_bootstrap_plan(&info, &SiteReplicationIamCredentials::default()).expect("bootstrap plan should build");
assert!(!plan.bucket_items.iter().any(|item| item.r#type == "lc-config"));
}
@@ -1748,7 +1865,8 @@ fn test_site_replication_bootstrap_plan_emits_timestamped_lifecycle_delete() {
},
);
let plan = site_replication_bootstrap_plan(&info).expect("bootstrap plan should build");
let plan =
site_replication_bootstrap_plan(&info, &SiteReplicationIamCredentials::default()).expect("bootstrap plan should build");
let item = plan
.bucket_items
@@ -1935,8 +2053,8 @@ fn test_site_replication_repair_preflight_token_is_deterministic_for_equal_state
},
);
let plan_a = site_replication_bootstrap_plan(&info).expect("first plan");
let plan_b = site_replication_bootstrap_plan(&info).expect("second plan");
let plan_a = site_replication_bootstrap_plan(&info, &SiteReplicationIamCredentials::default()).expect("first plan");
let plan_b = site_replication_bootstrap_plan(&info, &SiteReplicationIamCredentials::default()).expect("second plan");
let token_a = site_replication_repair_preflight_token(&state, &plan_a, b"test-signing-key").expect("first token");
let token_b = site_replication_repair_preflight_token(&state, &plan_b, b"test-signing-key").expect("second token");
@@ -3219,3 +3337,345 @@ fn test_reconcile_adds_missing_peer_rules_to_existing_config() {
assert!(rule_ids.contains(&"site-repl-dep-b"));
assert!(rule_ids.contains(&"site-repl-dep-c"));
}
/// backlog#2289: the IAM snapshot (retry resend, repair, site-add bootstrap)
/// used to be built from `list_users`, whose `UserInfo` never carries a
/// secret key, so the plan dropped every user and a status change or secret
/// rotation committed while a peer was unreachable never reached it. The
/// credentials now come from a separate store read; SRInfo stays secret-free.
#[test]
fn test_bootstrap_plan_carries_users_from_the_credential_snapshot() {
let mut info = SRInfo::default();
// Exactly what `list_users` builds: status, policy, updated_at — never secret_key.
info.user_info_map.insert(
"alice".to_string(),
rustfs_madmin::UserInfo {
status: rustfs_madmin::AccountStatus::Disabled,
policy_name: Some("readwrite".to_string()),
updated_at: Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp")),
..Default::default()
},
);
info.user_info_map.insert(
"external-idp-user".to_string(),
rustfs_madmin::UserInfo {
status: rustfs_madmin::AccountStatus::Enabled,
..Default::default()
},
);
let user_updated_at = OffsetDateTime::from_unix_timestamp(1_700_000_500).expect("timestamp");
let mut credentials = SiteReplicationIamCredentials::default();
credentials.users.insert(
"alice".to_string(),
SiteReplicationUserCredential {
secret_key: "alice-secret".to_string(),
status: rustfs_madmin::AccountStatus::Disabled,
updated_at: Some(user_updated_at),
},
);
let plan = site_replication_bootstrap_plan(&info, &credentials).expect("bootstrap plan should build");
let users: Vec<_> = plan.iam_items.iter().filter(|item| item.r#type == "iam-user").collect();
assert_eq!(users.len(), 1, "only the user with a credential travels: {:?}", plan.iam_items);
let alice = users[0].iam_user.as_ref().expect("iam user body");
assert_eq!(alice.access_key, "alice");
let req = alice.user_req.as_ref().expect("user request");
assert_eq!(req.secret_key, "alice-secret");
assert_eq!(req.status, rustfs_madmin::AccountStatus::Disabled);
assert_eq!(req.policy.as_deref(), Some("readwrite"));
// the user record's own axis, not the policy-mapping time list_users reports
assert_eq!(users[0].updated_at, Some(user_updated_at));
}
fn service_account_snapshot(access_key: &str, parent: &str, status: &str) -> SiteReplicationServiceAccountSnapshot {
SiteReplicationServiceAccountSnapshot {
create: rustfs_madmin::SRSvcAccCreate {
parent: parent.to_string(),
access_key: access_key.to_string(),
secret_key: format!("{access_key}-secret"),
groups: Vec::new(),
claims: HashMap::new(),
session_policy: SRSessionPolicy::default(),
status: status.to_string(),
name: String::new(),
description: String::new(),
expiration: None,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
},
envelope: None,
updated_at: Some(OffsetDateTime::from_unix_timestamp(1_700_000_600).expect("timestamp")),
}
}
/// backlog#2289: service accounts were absent from every snapshot (the
/// listing filters them). They now travel as the create item the live hook
/// emits — after their parents — carrying secret and status.
#[test]
fn test_bootstrap_plan_emits_service_accounts_after_their_parents() {
let mut info = SRInfo::default();
info.user_info_map
.insert("alice".to_string(), rustfs_madmin::UserInfo::default());
let mut credentials = SiteReplicationIamCredentials::default();
credentials.users.insert(
"alice".to_string(),
SiteReplicationUserCredential {
secret_key: "alice-secret".to_string(),
status: rustfs_madmin::AccountStatus::Enabled,
updated_at: None,
},
);
credentials
.service_accounts
.push(service_account_snapshot("alice-svc", "alice", "off"));
let plan = site_replication_bootstrap_plan(&info, &credentials).expect("bootstrap plan should build");
let types: Vec<_> = plan.iam_items.iter().map(|item| item.r#type.as_str()).collect();
assert_eq!(types, vec!["iam-user", "service-account"]);
let change = plan.iam_items[1].svc_acc_change.as_ref().expect("service account change");
let create = change.create.as_ref().expect("create body");
assert_eq!((create.access_key.as_str(), create.parent.as_str()), ("alice-svc", "alice"));
assert_eq!(create.secret_key, "alice-svc-secret");
assert_eq!(create.status, "off", "a disabled account must arrive disabled");
assert!(change.delete.is_none() && change.update.is_none());
}
/// A service account present in the previous snapshot but gone from the
/// fresh one is replayed as an explicit delete, like the other IAM kinds.
#[test]
fn test_retry_snapshot_tombstones_removed_service_accounts() {
let observed_at = OffsetDateTime::from_unix_timestamp(1_700_001_000).expect("timestamp");
let mut info = SRInfo::default();
info.user_info_map
.insert("alice".to_string(), rustfs_madmin::UserInfo::default());
let mut credentials = SiteReplicationIamCredentials::default();
credentials.users.insert(
"alice".to_string(),
SiteReplicationUserCredential {
secret_key: "alice-secret".to_string(),
status: rustfs_madmin::AccountStatus::Enabled,
updated_at: None,
},
);
let mut with_account = credentials.clone();
with_account
.service_accounts
.push(service_account_snapshot("alice-svc", "alice", "on"));
let previous = site_replication_bootstrap_plan(&info, &with_account).expect("previous plan");
let fresh = site_replication_bootstrap_plan(&info, &credentials).expect("fresh plan");
let replay = RetrySnapshot::replay_after_change(
&RetrySnapshot::Iam(previous.iam_items),
&RetrySnapshot::Iam(fresh.iam_items),
observed_at,
);
let RetrySnapshot::Iam(items) = replay else {
panic!("IAM snapshot expected");
};
let tombstone = items
.iter()
.find(|item| item.r#type == "service-account")
.expect("service account tombstone");
let change = tombstone.svc_acc_change.as_ref().expect("change");
assert_eq!(change.delete.as_ref().map(|delete| delete.access_key.as_str()), Some("alice-svc"));
assert!(change.create.is_none());
assert_eq!(tombstone.updated_at, Some(observed_at));
}
/// Spawns a one-shot HTTP peer that answers 200 and flips the returned flag
/// once a request head has arrived.
async fn spawn_reached_probe_peer() -> (String, Arc<AtomicBool>, tokio::task::JoinHandle<()>) {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind healthy peer");
let endpoint = format!("http://{}", listener.local_addr().expect("healthy peer address"));
let reached = Arc::new(AtomicBool::new(false));
let reached_by_server = reached.clone();
let server = tokio::spawn(async move {
let Ok((mut stream, _)) = listener.accept().await else {
return;
};
let mut request = Vec::new();
let mut buffer = [0_u8; 1024];
loop {
let Ok(read) = stream.read(&mut buffer).await else {
return;
};
if read == 0 {
return;
}
request.extend_from_slice(&buffer[..read]);
if request.windows(4).any(|window| window == b"\r\n\r\n") {
break;
}
}
reached_by_server.store(true, Ordering::SeqCst);
let _ = stream
.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 2\r\nconnection: close\r\n\r\nok")
.await;
});
(endpoint, reached, server)
}
/// Three-peer runtime whose local peer is `local`; BTreeMap order visits the
/// failing peer `b` before the healthy peer `c`.
fn broadcast_runtime_with_failing_peer_before_healthy(failing_endpoint: &str, healthy_endpoint: &str) -> SiteReplicationRuntime {
let local_peer = PeerInfo {
deployment_id: "local".to_string(),
..peer("local", "http://127.0.0.1:9")
};
let mut state = SiteReplicationState {
name: "local".to_string(),
service_account_access_key: "site-replicator-0".to_string(),
..Default::default()
};
state.peers.insert("local".to_string(), local_peer.clone());
state.peers.insert(
"b".to_string(),
PeerInfo {
deployment_id: "b".to_string(),
..peer("b", failing_endpoint)
},
);
state.peers.insert(
"c".to_string(),
PeerInfo {
deployment_id: "c".to_string(),
..peer("c", healthy_endpoint)
},
);
SiteReplicationRuntime {
state,
local_peer,
service_account_secret_key: "site-replicator-secret".to_string(),
}
}
const BROADCAST_PROBE_DELETE_BUCKET_PATH: &str =
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=delete-bucket";
/// The generic JSON broadcast (bucket make/delete, bucket-meta hook, bucket
/// ops) attempts every remote peer: a peer whose request fails must not stop
/// delivery to the peers that follow it in deployment-id order, and the
/// failure is still reported to the caller (backlog#2293).
#[tokio::test]
#[serial]
async fn test_broadcast_json_reaches_healthy_peers_after_a_failed_peer() {
// Peer "b": nothing listens on the port, so the connect is refused.
let refused = TcpListener::bind("127.0.0.1:0").await.expect("bind refused-peer probe");
let refused_endpoint = format!("http://{}", refused.local_addr().expect("refused-peer address"));
drop(refused);
let (healthy_endpoint, reached, server) = spawn_reached_probe_peer().await;
let runtime = broadcast_runtime_with_failing_peer_before_healthy(&refused_endpoint, &healthy_endpoint);
let result = temp_env::async_with_vars([(ALLOW_LOOPBACK_REPLICATION_TARGET_ENV, Some("true"))], async {
broadcast_site_replication_json_with_runtime(&runtime, BROADCAST_PROBE_DELETE_BUCKET_PATH, &serde_json::json!({})).await
})
.await;
let err = result.expect_err("peer b refuses connections, the broadcast must report it");
assert!(
reached.load(Ordering::SeqCst),
"peer c never received the broadcast once peer b failed: {err}"
);
server.abort();
}
/// Same guarantee when the failing peer never gets a transport: an endpoint
/// that `PeerTransport::for_runtime_peer` rejects must be skipped past (and
/// reported), not abort the broadcast before the healthy peers (backlog#2293).
#[tokio::test]
#[serial]
async fn test_broadcast_json_reaches_healthy_peers_after_a_peer_without_transport() {
// Peer "b": a scheme the peer connection validator refuses outright.
let forbidden_endpoint = "ftp://peer-b.example.com";
let (healthy_endpoint, reached, server) = spawn_reached_probe_peer().await;
let runtime = broadcast_runtime_with_failing_peer_before_healthy(forbidden_endpoint, &healthy_endpoint);
let result = temp_env::async_with_vars([(ALLOW_LOOPBACK_REPLICATION_TARGET_ENV, Some("true"))], async {
broadcast_site_replication_json_with_runtime(&runtime, BROADCAST_PROBE_DELETE_BUCKET_PATH, &serde_json::json!({})).await
})
.await;
let err = result.expect_err("peer b has no usable transport, the broadcast must report it");
assert!(
err.to_string().contains("invalid persisted site replication peer"),
"the reported error must be peer b's transport failure: {err}"
);
assert!(
reached.load(Ordering::SeqCst),
"peer c never received the broadcast once peer b failed to get a transport: {err}"
);
server.abort();
}
fn service_account_item_with_claims(order: &[&str]) -> SRIAMItem {
let mut claims = HashMap::new();
for key in order {
claims.insert((*key).to_string(), serde_json::json!(format!("value-of-{key}")));
}
SRIAMItem {
r#type: "service-account".to_string(),
svc_acc_change: Some(SRSvcAccChange {
create: Some(rustfs_madmin::SRSvcAccCreate {
parent: "alice".to_string(),
access_key: "alice-svc".to_string(),
secret_key: "alice-svc-secret".to_string(),
groups: Vec::new(),
claims,
session_policy: SRSessionPolicy::default(),
status: "on".to_string(),
name: String::new(),
description: String::new(),
expiration: None,
api_version: Some(SITE_REPL_API_VERSION.to_string()),
}),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
}),
updated_at: Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp")),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
}
}
/// The repair preflight token and the retry-snapshot fingerprint hash the
/// serialized items. Service-account claims live in a `HashMap`, whose
/// iteration order differs between instances, so the hash must not depend on
/// it (the real-VM repair returned 412 "preflight is stale" between dry-run
/// and execute once snapshots carried service accounts).
#[test]
fn test_repair_task_id_and_retry_fingerprint_ignore_claim_map_order() {
let forward = service_account_item_with_claims(&["accessKey", "exp", "parent", "sa-policy", "sub", "tenant"]);
let backward = service_account_item_with_claims(&["tenant", "sub", "sa-policy", "parent", "exp", "accessKey"]);
let canonical = canonical_json_vec(&forward).expect("canonical json");
let text = String::from_utf8(canonical).expect("utf8");
let positions: Vec<usize> = [
"\"accessKey\"",
"\"exp\"",
"\"parent\"",
"\"sa-policy\"",
"\"sub\"",
"\"tenant\"",
]
.iter()
.map(|key| text.find(key).expect("claim key present"))
.collect();
assert!(
positions.windows(2).all(|pair| pair[0] < pair[1]),
"claim keys must serialize sorted: {text}"
);
assert_eq!(
SiteReplicationRepairTask::Iam(&forward).id().expect("id"),
SiteReplicationRepairTask::Iam(&backward).id().expect("id"),
"identical items must yield the same repair task id regardless of claim map order"
);
assert_eq!(
RetrySnapshot::Iam(vec![forward]).fingerprint().expect("fingerprint"),
RetrySnapshot::Iam(vec![backward]).fingerprint().expect("fingerprint"),
"identical snapshots must fingerprint equal regardless of claim map order"
);
}
+24 -6
View File
@@ -876,6 +876,14 @@ pub(crate) async fn broadcast_site_replication_json<T: Serialize>(path: &str, bo
broadcast_site_replication_json_with_runtime(&runtime, path, body).await
}
/// PUT `body` to `path` on every remote peer of the runtime.
///
/// Every peer is attempted: one peer's failure — transport construction
/// included — must not skip the peers that follow it in deployment-id order,
/// or they silently miss the change with no retry record (backlog#2293). A
/// success settles the peer/path's queued retry event, a failure enqueues one
/// under the request `path` (so the drain classifies it as today), and the
/// first error is returned once all peers were attempted.
pub(crate) async fn broadcast_site_replication_json_with_runtime<T: Serialize>(
runtime: &SiteReplicationRuntime,
path: &str,
@@ -883,20 +891,30 @@ pub(crate) async fn broadcast_site_replication_json_with_runtime<T: Serialize>(
) -> S3Result<()> {
let state = &runtime.state;
let local_peer = &runtime.local_peer;
let mut first_error: Option<S3Error> = None;
for peer in state.peers.values() {
if peer.deployment_id == local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) {
continue;
}
let transport = PeerTransport::for_runtime_peer(peer).await?;
PeerAdminRequest::put(&transport.connection, path, &state.service_account_access_key)
.with_client(&transport.client)
.send_with_retry_event(peer, &runtime.service_account_secret_key, body)
.await?;
let sent = match PeerTransport::for_runtime_peer(peer).await {
Ok(transport) => PeerAdminRequest::put(&transport.connection, path, &state.service_account_access_key)
.with_client(&transport.client)
.send_with_retry_event(peer, &runtime.service_account_secret_key, body)
.await
.map(|_| ()),
Err(err) => {
enqueue_site_replication_retry_event(peer, path, &err).await;
Err(err)
}
};
if let Err(err) = sent {
first_error.get_or_insert(err);
}
}
Ok(())
first_error.map_or(Ok(()), Err)
}
pub(crate) fn parse_endpoint_refresh_status(peer: &PeerInfo, body: &[u8]) -> S3Result<()> {