commit 45320c7d54da966c2c46f2ff933fc4a70c357336
parent dfec133745551bc6046552792a5b0fb1f07e703d
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Tue, 25 Aug 2026 13:03:07 +0200
Invalidate stale sync work on chain reset
Diffstat:
4 files changed, 88 insertions(+), 24 deletions(-)
diff --git a/src/adapters/http/actions.rs b/src/adapters/http/actions.rs
@@ -321,6 +321,7 @@ pub(super) async fn reset_local_chain(state: &HttpState, confirmation: &str) ->
{
let mut node = state.node.lock().await;
+ state.gossip.invalidate_sync_progress();
node.reset_chain_to_setup_placeholder();
}
clear_chain(&state.chain_store).await?;
diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs
@@ -99,12 +99,17 @@ pub struct GossipNetwork {
pub(super) struct SyncProgressGuard {
network: GossipNetwork,
id: u64,
+ generation: u64,
}
impl SyncProgressGuard {
pub(super) fn id(&self) -> u64 {
self.id
}
+
+ pub(super) fn is_current(&self) -> bool {
+ self.network.sync_generation_is_current(self.generation)
+ }
}
impl Drop for SyncProgressGuard {
@@ -136,6 +141,7 @@ pub struct SyncProgress {
#[derive(Default)]
struct SyncProgressState {
next_id: u64,
+ generation: u64,
active: BTreeMap<u64, SyncProgress>,
}
diff --git a/src/adapters/p2p/network.rs b/src/adapters/p2p/network.rs
@@ -142,9 +142,32 @@ impl GossipNetwork {
super::SyncProgressGuard {
network: self.clone(),
id,
+ generation: state.generation,
}
}
+ pub(crate) fn invalidate_sync_progress(&self) {
+ let mut state = self
+ .inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned");
+ state.generation = state.generation.wrapping_add(1);
+ state.active.clear();
+ }
+
+ pub(super) fn sync_generation(&self) -> u64 {
+ self.inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned")
+ .generation
+ }
+
+ pub(super) fn sync_generation_is_current(&self, generation: u64) -> bool {
+ self.sync_generation() == generation
+ }
+
pub(super) fn update_sync_progress(&self, id: u64, validated_height: u64) {
let mut state = self
.inner
@@ -357,6 +380,33 @@ mod tests {
}
#[tokio::test]
+ async fn invalidating_sync_progress_hides_and_rejects_stale_validation() {
+ let alice = Wallet::from_seed("invalidated-sync-progress-alice");
+ let allocations = allocations(std::slice::from_ref(&alice), 1_000);
+ let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
+ let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
+ let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None);
+
+ let stale = network.begin_sync_progress(20, 90);
+ network.update_sync_progress(stale.id(), 41);
+ assert!(stale.is_current());
+
+ network.invalidate_sync_progress();
+
+ assert_eq!(network.sync_progress(), None);
+ assert!(!stale.is_current());
+ network.update_sync_progress(stale.id(), 42);
+ assert_eq!(network.sync_progress(), None);
+
+ let fresh = network.begin_sync_progress(0, 90);
+ assert!(fresh.is_current());
+ assert_eq!(network.sync_progress().unwrap().validated_height, 0);
+
+ drop(stale);
+ assert_eq!(network.sync_progress().unwrap().validated_height, 0);
+ }
+
+ #[tokio::test]
async fn peer_exchange_does_not_advertise_self_when_outbound_only() {
let alice = Wallet::from_seed("px-private-alice");
let allocations = allocations(std::slice::from_ref(&alice), 1_000);
diff --git a/src/adapters/p2p/process.rs b/src/adapters/p2p/process.rs
@@ -161,16 +161,20 @@ pub(super) async fn process_envelope(
}
GossipEnvelope::Blocks { blocks } => {
let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
- let local_ledger = network.inner.node.lock().await.clone_ledger();
- let start_height = blocks
- .first()
- .map(|block| block.height.saturating_sub(1))
- .unwrap_or_else(|| local_ledger.height());
- let target_height = blocks
- .last()
- .map(|block| block.height)
- .unwrap_or(start_height);
- let progress_guard = network.begin_sync_progress(start_height, target_height);
+ let (local_ledger, progress_guard) = {
+ let node = network.inner.node.lock().await;
+ let local_ledger = node.clone_ledger();
+ let start_height = blocks
+ .first()
+ .map(|block| block.height.saturating_sub(1))
+ .unwrap_or_else(|| local_ledger.height());
+ let target_height = blocks
+ .last()
+ .map(|block| block.height)
+ .unwrap_or(start_height);
+ let progress_guard = network.begin_sync_progress(start_height, target_height);
+ (local_ledger, progress_guard)
+ };
let progress_id = progress_guard.id();
let progress_network = network.clone();
let result = match validate_blocks_extension(
@@ -181,13 +185,14 @@ pub(super) async fn process_envelope(
)
.await
{
- Ok(ledger) => network
- .inner
- .node
- .lock()
- .await
- .import_verified_ledger(ledger)
- .map(|_| ()),
+ Ok(ledger) => {
+ let mut node = network.inner.node.lock().await;
+ if progress_guard.is_current() {
+ node.import_verified_ledger(ledger).map(|_| ())
+ } else {
+ Ok(())
+ }
+ }
Err(error) => Err(error),
};
drop(progress_guard);
@@ -205,15 +210,17 @@ pub(super) async fn process_envelope(
network.forward_outbox().await;
}
GossipEnvelope::ChainBootstrap(bootstrap) => {
+ let sync_generation = network.sync_generation();
let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
let result = match validate_chain_bootstrap(bootstrap, adjusted_time_ms).await {
- Ok(ledger) => network
- .inner
- .node
- .lock()
- .await
- .import_verified_ledger(ledger)
- .map(|_| ()),
+ Ok(ledger) => {
+ let mut node = network.inner.node.lock().await;
+ if network.sync_generation_is_current(sync_generation) {
+ node.import_verified_ledger(ledger).map(|_| ())
+ } else {
+ Ok(())
+ }
+ }
Err(error) => Err(error),
};
record_rejected_chain_payload(