Skip to content

Commit 584a1bc

Browse files
authored
🐛 fix(ha): fence cross-datacenter transfers (#1859)
* 🐛 fix(ha): fence cross-DC transfers Pending placements survived a restart without a durable way to distinguish an abandoned transfer from a live copy. That blocked recovery or risked two owners publishing the same placement. Bind each transfer to the consensus term and a monotonic attempt token. A newer term can reclaim pending work, while stale checkpoints and terminal writes fail their ownership check. * 🧪 test(ha): set placement attempt
1 parent 0da6835 commit 584a1bc

22 files changed

Lines changed: 414 additions & 71 deletions

crates/peryx-driver/tests/unit/availability/tests.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,7 @@ fn blob_placement_view_projects_and_orders_states() {
121121
},
122122
state: state_value,
123123
fence: 1,
124+
transfer_attempt: 1,
124125
generation: 1,
125126
updated_at_unix: updated_at,
126127
};

crates/peryx-ha-distributed/src/copy_planning.rs

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,11 @@ pub struct CrossDcCopy {
2525
pub fence: NonZeroU64,
2626
}
2727

28-
pub fn copy_backlog_entry(records: &[BlobPlacementRecord], local_dc: &DataCenterId) -> Option<CopyBacklogEntry> {
28+
pub fn copy_backlog_entry(
29+
records: &[BlobPlacementRecord],
30+
local_dc: &DataCenterId,
31+
fence: NonZeroU64,
32+
) -> Option<CopyBacklogEntry> {
2933
let mut local_settled = false;
3034
let mut sources = Vec::new();
3135
for record in records {
@@ -36,11 +40,10 @@ pub fn copy_backlog_entry(records: &[BlobPlacementRecord], local_dc: &DataCenter
3640
generation: record.generation,
3741
size,
3842
}),
39-
BlobPlacementState::Verified { .. } | BlobPlacementState::Pending | BlobPlacementState::Revoked
40-
if is_local =>
41-
{
43+
BlobPlacementState::Verified { .. } | BlobPlacementState::Revoked if is_local => {
4244
local_settled = true;
4345
}
46+
BlobPlacementState::Pending if is_local && record.fence >= fence.get() => local_settled = true,
4447
_ => {}
4548
}
4649
}

crates/peryx-ha-distributed/src/copy_runtime.rs

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ use crate::{BlobTransport, CopyError, HttpBlobTransport, TransferLimits, Transpo
1313
use peryx_core::Clock;
1414
use peryx_ha::{
1515
AvailabilityTaskError, AvailabilityTaskReport, BackendId, BlobPlacementFailure, BlobPlacementKey,
16-
BlobPlacementTransition, DataCenterId,
16+
BlobPlacementOutcome, BlobPlacementTransition, DataCenterId,
1717
};
1818
use peryx_storage::blob::{BlobStore, Digest};
1919
use peryx_storage::meta::MetaStore;
@@ -95,7 +95,7 @@ impl CrossDcBlobCopier {
9595
.scan_blob_placement_groups(cursor.as_deref(), batch)
9696
.map_err(|error| task_error("copy_backlog_scan", error))?;
9797
planned.extend(page.groups.iter().filter_map(|records| {
98-
copy_backlog_entry(records, &self.local_dc)
98+
copy_backlog_entry(records, &self.local_dc, fence)
9999
.map(|entry| plan_cross_dc_copy(&entry, &self.local_dc, &self.backend, fence))
100100
}));
101101
match page.next_cursor {
@@ -112,27 +112,29 @@ impl CrossDcBlobCopier {
112112
tracing::warn!(source_dc, "cross-datacenter copy has no reachable source peer");
113113
return false;
114114
};
115-
if !record(
115+
let Some(BlobPlacementOutcome::Applied(staged)) = record(
116116
meta,
117117
&copy.target,
118118
&BlobPlacementTransition::Stage,
119119
copy.fence.get(),
120120
clock,
121-
) {
121+
) else {
122122
return false;
123-
}
123+
};
124124
let digest = Digest::from_hex(copy.target.digest.sha256()).expect("artifact digests are validated SHA-256");
125125
let outcome = copy_blob_to_target(transport, &self.store, &digest).await;
126126
let transition = match &outcome {
127127
Ok(()) => BlobPlacementTransition::Verify {
128+
attempt: staged.transfer_attempt,
128129
observed: copy.target.digest.clone(),
129130
size: copy.size,
130131
},
131132
Err(error) => BlobPlacementTransition::Fail {
133+
attempt: staged.transfer_attempt,
132134
class: failure_class(error),
133135
},
134136
};
135-
let recorded = record(meta, &copy.target, &transition, copy.fence.get(), clock);
137+
let recorded = record(meta, &copy.target, &transition, copy.fence.get(), clock).is_some();
136138
outcome.is_ok() && recorded
137139
}
138140
}
@@ -173,12 +175,12 @@ fn record(
173175
transition: &BlobPlacementTransition,
174176
fence: u64,
175177
clock: &Clock,
176-
) -> bool {
178+
) -> Option<BlobPlacementOutcome> {
177179
match apply_blob_placement(meta, key, transition, fence, (clock)()) {
178-
Ok(_) => true,
180+
Ok(outcome) => Some(outcome),
179181
Err(error) => {
180182
tracing::warn!(%error, ?transition, "cross-datacenter copy could not record a placement");
181-
false
183+
None
182184
}
183185
}
184186
}

crates/peryx-ha-distributed/src/placement_policy.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -96,11 +96,12 @@ pub fn record_local_placement(
9696
{
9797
return Ok(BlobPlacementOutcome::Unchanged(record));
9898
}
99-
apply_blob_placement(meta, &key, &BlobPlacementTransition::Stage, fence, now)?;
99+
let staged = apply_blob_placement(meta, &key, &BlobPlacementTransition::Stage, fence, now)?;
100100
apply_blob_placement(
101101
meta,
102102
&key,
103103
&BlobPlacementTransition::Verify {
104+
attempt: staged.record().transfer_attempt,
104105
observed: digest.clone(),
105106
size,
106107
},

crates/peryx-ha-distributed/src/placement_runtime.rs

Lines changed: 3 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@ use std::num::NonZeroUsize;
66

77
use peryx_core::Clock;
88
use peryx_ha::{
9-
AvailabilityTaskError, AvailabilityTaskReport, BlobPlacementFailure, BlobPlacementKey, BlobPlacementState,
10-
BlobPlacementTransition, DataCenterId, ReclamationState, ReclamationStore,
9+
AvailabilityTaskError, AvailabilityTaskReport, BlobPlacementKey, BlobPlacementState, BlobPlacementTransition,
10+
DataCenterId, ReclamationState, ReclamationStore,
1111
};
1212
use peryx_identity::ArtifactDigest;
1313
use peryx_storage::blob::{BlobErrorKind, BlobStore, Digest};
@@ -97,15 +97,7 @@ impl FilesystemPlacementReconciler {
9797
return false;
9898
}
9999
// Stop routing to corrupt bytes before clearing the path for a replacement.
100-
if !record_transition(
101-
meta,
102-
key,
103-
&BlobPlacementTransition::Fail {
104-
class: BlobPlacementFailure::DigestMismatch,
105-
},
106-
fence,
107-
clock,
108-
) {
100+
if !record_transition(meta, key, &BlobPlacementTransition::Invalidate, fence, clock) {
109101
return false;
110102
}
111103
if let Err(error) = self.store.remove(&digest) {

crates/peryx-ha-distributed/tests/unit/blob_placements_http_tests.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ async fn build_app(corrupt: bool) -> (tempfile::TempDir, Arc<AppState>) {
5252
drop(authorization);
5353
drop(users);
5454
let verify = BlobPlacementTransition::Verify {
55+
attempt: 1,
5556
observed: digest(),
5657
size: 4096,
5758
};
@@ -63,6 +64,7 @@ async fn build_app(corrupt: bool) -> (tempfile::TempDir, Arc<AppState>) {
6364
&[
6465
BlobPlacementTransition::Stage,
6566
BlobPlacementTransition::Fail {
67+
attempt: 1,
6668
class: BlobPlacementFailure::SourceUnavailable,
6769
},
6870
],

crates/peryx-ha-distributed/tests/unit/blob_plane_tests.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,7 @@ fn seed_verified_placement_on(meta: &MetaStore, digest: &Digest, dc: &str, backe
7878
meta,
7979
&key,
8080
&BlobPlacementTransition::Verify {
81+
attempt: 1,
8182
observed: artifact,
8283
size,
8384
},

crates/peryx-ha-distributed/tests/unit/copy_planning_tests.rs

Lines changed: 28 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ fn record(
3030
key: key(suffix, data_center, location),
3131
state,
3232
fence: 1,
33+
transfer_attempt: 1,
3334
generation,
3435
updated_at_unix: 0,
3536
}
@@ -43,6 +44,7 @@ fn test_backlog_selects_verified_remote_sources() {
4344
record(1, "south", "south/01", 4, BlobPlacementState::Verified { size: 20 }),
4445
],
4546
&DataCenterId::new("west").unwrap(),
47+
NonZeroU64::new(1).unwrap(),
4648
)
4749
.unwrap();
4850

@@ -66,6 +68,7 @@ fn test_backlog_breaks_generation_ties_by_ledger_order() {
6668
record(1, "south", "south/01", 4, BlobPlacementState::Verified { size: 20 }),
6769
],
6870
&DataCenterId::new("west").unwrap(),
71+
NonZeroU64::new(1).unwrap(),
6972
)
7073
.unwrap();
7174

@@ -94,6 +97,7 @@ fn test_backlog_skips_settled_local_placements() {
9497
record(1, "west", "west/01", 1, state),
9598
],
9699
&local,
100+
NonZeroU64::new(1).unwrap(),
97101
)
98102
.is_none()
99103
);
@@ -116,6 +120,21 @@ fn test_backlog_retries_failed_local_placements() {
116120
),
117121
],
118122
&DataCenterId::new("west").unwrap(),
123+
NonZeroU64::new(1).unwrap(),
124+
);
125+
126+
assert!(entry.is_some());
127+
}
128+
129+
#[test]
130+
fn test_backlog_retries_pending_after_an_ownership_transition() {
131+
let entry = copy_backlog_entry(
132+
&[
133+
record(1, "east", "east/01", 1, BlobPlacementState::Verified { size: 10 }),
134+
record(1, "west", "west/01", 1, BlobPlacementState::Pending),
135+
],
136+
&DataCenterId::new("west").unwrap(),
137+
NonZeroU64::new(2).unwrap(),
119138
);
120139

121140
assert!(entry.is_some());
@@ -125,6 +144,13 @@ fn test_backlog_retries_failed_local_placements() {
125144
fn test_backlog_requires_records_and_a_verified_source() {
126145
let local = DataCenterId::new("west").unwrap();
127146

128-
assert!(copy_backlog_entry(&[], &local).is_none());
129-
assert!(copy_backlog_entry(&[record(1, "east", "east/01", 1, BlobPlacementState::Pending)], &local).is_none());
147+
assert!(copy_backlog_entry(&[], &local, NonZeroU64::new(1).unwrap()).is_none());
148+
assert!(
149+
copy_backlog_entry(
150+
&[record(1, "east", "east/01", 1, BlobPlacementState::Pending)],
151+
&local,
152+
NonZeroU64::new(1).unwrap(),
153+
)
154+
.is_none()
155+
);
130156
}

0 commit comments

Comments
 (0)