commit 925c2887f47d16746138aed90a07d599ea269d2e
parent 63999847d979e96a008c45b712ceb8d93b842e8c
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Mon, 24 Aug 2026 12:18:26 +0200
Report incremental block sync progress
Diffstat:
11 files changed, 259 insertions(+), 22 deletions(-)
diff --git a/src/adapters/http/api.rs b/src/adapters/http/api.rs
@@ -409,6 +409,7 @@ fn metrics_leaderboard_entry(
pub(super) async fn api_network_health(
State(state): State<HttpState>,
) -> Json<NetworkHealthResponse> {
+ let sync_progress = state.gossip.sync_progress();
let (local, mempool) = {
let node = state.node.lock().await;
let status = node.status();
@@ -426,6 +427,9 @@ pub(super) async fn api_network_health(
height: status.chain.height,
tip_hash: status.chain.tip_hash,
tip_timestamp_ms,
+ sync_start_height: sync_progress.map(|progress| progress.start_height),
+ sync_validated_height: sync_progress.map(|progress| progress.validated_height),
+ sync_target_height: sync_progress.map(|progress| progress.target_height),
pending_transactions: mempool.total(),
last_finalizer_mode: tip.as_ref().map(|block| match block.finalizer_mode {
crate::domain::FinalizerMode::Ticket => "ticket".to_string(),
diff --git a/src/adapters/http/metrics.rs b/src/adapters/http/metrics.rs
@@ -169,7 +169,10 @@ pub(super) fn network_health_at(
.tip_timestamp_ms
.map(|tip_timestamp_ms| now_ms.saturating_sub(tip_timestamp_ms));
let remote_best_height = peers.iter().filter_map(|peer| peer.last_known_height).max();
- let best_known_height = remote_best_height.unwrap_or(local_height).max(local_height);
+ let best_known_height = remote_best_height
+ .unwrap_or(local_height)
+ .max(local.sync_target_height.unwrap_or(local_height))
+ .max(local_height);
let healthy_heights = peers
.iter()
.filter(|peer| peer.last_error.is_none())
@@ -229,8 +232,13 @@ pub(super) fn network_health_at(
.rejected_blocks
.saturating_add(local.rejected_block_batches)
.saturating_add(local.rejected_snapshots);
+ let actively_syncing = local
+ .sync_target_height
+ .is_some_and(|target_height| target_height > local_height);
- let state = if peers.is_empty() {
+ let state = if actively_syncing {
+ "syncing"
+ } else if peers.is_empty() {
"isolated"
} else if banned_peers > 0 && healthy_peers == 0 {
"banned"
@@ -254,6 +262,9 @@ pub(super) fn network_health_at(
local_tip_hash: local.tip_hash,
last_block_age_ms,
best_known_height,
+ sync_start_height: local.sync_start_height,
+ sync_validated_height: local.sync_validated_height,
+ sync_target_height: local.sync_target_height,
shared_height,
lag_blocks,
outbound_peers,
@@ -352,6 +363,9 @@ mod tests {
height: 42,
tip_hash: "tip-hash".to_string(),
tip_timestamp_ms: Some(1_000),
+ sync_start_height: Some(42),
+ sync_validated_height: Some(47),
+ sync_target_height: Some(60),
pending_transactions: 3,
last_finalizer_mode: Some("ticket".to_string()),
last_finalizer_rank: Some(1),
@@ -377,6 +391,11 @@ mod tests {
);
assert_eq!(health.local_height, 42);
+ assert_eq!(health.best_known_height, 60);
+ assert_eq!(health.sync_start_height, Some(42));
+ assert_eq!(health.sync_validated_height, Some(47));
+ assert_eq!(health.sync_target_height, Some(60));
+ assert_eq!(health.state, "syncing");
assert_eq!(health.local_tip_hash, "tip-hash");
assert_eq!(health.last_block_age_ms, Some(1_500));
assert_eq!(health.last_finalizer_mode.as_deref(), Some("ticket"));
diff --git a/src/adapters/http/types.rs b/src/adapters/http/types.rs
@@ -30,6 +30,9 @@ pub(super) struct NetworkHealthResponse {
pub(super) local_tip_hash: String,
pub(super) last_block_age_ms: Option<u64>,
pub(super) best_known_height: u64,
+ pub(super) sync_start_height: Option<u64>,
+ pub(super) sync_validated_height: Option<u64>,
+ pub(super) sync_target_height: Option<u64>,
pub(super) shared_height: u64,
pub(super) lag_blocks: u64,
pub(super) outbound_peers: usize,
@@ -74,6 +77,9 @@ pub(super) struct NetworkHealthLocalState {
pub(super) height: u64,
pub(super) tip_hash: String,
pub(super) tip_timestamp_ms: Option<u64>,
+ pub(super) sync_start_height: Option<u64>,
+ pub(super) sync_validated_height: Option<u64>,
+ pub(super) sync_target_height: Option<u64>,
pub(super) pending_transactions: usize,
pub(super) last_finalizer_mode: Option<String>,
pub(super) last_finalizer_rank: Option<u32>,
diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs
@@ -96,6 +96,23 @@ pub struct GossipNetwork {
inner: Arc<GossipNetworkInner>,
}
+pub(super) struct SyncProgressGuard {
+ network: GossipNetwork,
+ id: u64,
+}
+
+impl SyncProgressGuard {
+ pub(super) fn id(&self) -> u64 {
+ self.id
+ }
+}
+
+impl Drop for SyncProgressGuard {
+ fn drop(&mut self) {
+ self.network.finish_sync_progress(self.id);
+ }
+}
+
struct GossipNetworkInner {
node: SharedNode,
peers: SharedPeerBook,
@@ -106,6 +123,20 @@ struct GossipNetworkInner {
sessions: Mutex<BTreeMap<String, mpsc::Sender<OutboundBatch>>>,
inbound_limiter: Arc<StdMutex<InboundConnectionLimiter>>,
metrics: P2pMetricsCounters,
+ sync_progress: StdMutex<SyncProgressState>,
+}
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct SyncProgress {
+ pub start_height: u64,
+ pub validated_height: u64,
+ pub target_height: u64,
+}
+
+#[derive(Default)]
+struct SyncProgressState {
+ next_id: u64,
+ active: BTreeMap<u64, SyncProgress>,
}
#[cfg(test)]
diff --git a/src/adapters/p2p/fetch.rs b/src/adapters/p2p/fetch.rs
@@ -208,12 +208,14 @@ pub(super) async fn validate_blocks_extension(
mut ledger: Ledger,
blocks: Vec<Block>,
now_ms: u64,
+ on_progress: impl Fn(u64) + Send + 'static,
) -> Result<Ledger> {
if blocks.is_empty() {
return Ok(ledger);
}
tokio::task::spawn_blocking(move || {
+ let target_height = blocks.last().map(|block| block.height);
if blocks[0].prev_hash != ledger.tip_hash() {
let mut candidate = ledger.snapshot();
let ancestor = candidate
@@ -224,10 +226,15 @@ pub(super) async fn validate_blocks_extension(
candidate.blocks.truncate(ancestor + 1);
candidate.blocks.extend(blocks);
ledger.extend_from_snapshot_at(candidate, now_ms)?;
+ if let Some(target_height) = target_height {
+ on_progress(target_height);
+ }
return Ok(ledger);
}
for block in blocks {
+ let height = block.height;
ledger.apply_block_at(block, now_ms)?;
+ on_progress(height);
}
Ok(ledger)
})
diff --git a/src/adapters/p2p/network.rs b/src/adapters/p2p/network.rs
@@ -40,6 +40,7 @@ impl GossipNetwork {
),
inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())),
metrics: P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
@@ -107,6 +108,64 @@ impl GossipNetwork {
self.inner.metrics.snapshot()
}
+ pub fn sync_progress(&self) -> Option<super::SyncProgress> {
+ self.inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned")
+ .active
+ .values()
+ .copied()
+ .max_by_key(|progress| (progress.target_height, progress.validated_height))
+ }
+
+ pub(super) fn begin_sync_progress(
+ &self,
+ start_height: u64,
+ target_height: u64,
+ ) -> super::SyncProgressGuard {
+ let mut state = self
+ .inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned");
+ state.next_id = state.next_id.wrapping_add(1);
+ let id = state.next_id;
+ state.active.insert(
+ id,
+ super::SyncProgress {
+ start_height,
+ validated_height: start_height,
+ target_height: target_height.max(start_height),
+ },
+ );
+ super::SyncProgressGuard {
+ network: self.clone(),
+ id,
+ }
+ }
+
+ pub(super) fn update_sync_progress(&self, id: u64, validated_height: u64) {
+ let mut state = self
+ .inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned");
+ if let Some(progress) = state.active.get_mut(&id) {
+ progress.validated_height =
+ validated_height.clamp(progress.start_height, progress.target_height);
+ }
+ }
+
+ pub(super) fn finish_sync_progress(&self, id: u64) {
+ let mut state = self
+ .inner
+ .sync_progress
+ .lock()
+ .expect("sync progress mutex poisoned");
+ state.active.remove(&id);
+ }
+
pub(super) fn try_acquire_inbound_session(
&self,
ip: IpAddr,
@@ -269,6 +328,35 @@ mod tests {
use super::super::test_support::{allocations, gossip_network, node};
#[tokio::test]
+ async fn sync_progress_is_incremental_and_scoped_to_the_active_validation() {
+ let alice = Wallet::from_seed("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 first = network.begin_sync_progress(0, 90);
+ network.update_sync_progress(first.id(), 23);
+ assert_eq!(
+ network.sync_progress(),
+ Some(super::super::SyncProgress {
+ start_height: 0,
+ validated_height: 23,
+ target_height: 90,
+ })
+ );
+
+ let second = network.begin_sync_progress(23, 76);
+ network.update_sync_progress(second.id(), 41);
+ assert_eq!(network.sync_progress().unwrap().target_height, 90);
+ drop(first);
+ assert_eq!(network.sync_progress().unwrap().validated_height, 41);
+
+ drop(second);
+ assert_eq!(network.sync_progress(), None);
+ }
+
+ #[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
@@ -162,17 +162,35 @@ 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 result =
- match validate_blocks_extension(local_ledger, blocks, adjusted_time_ms).await {
- Ok(ledger) => network
- .inner
- .node
- .lock()
- .await
- .import_verified_ledger(ledger)
- .map(|_| ()),
- Err(error) => Err(error),
- };
+ 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 progress_id = progress_guard.id();
+ let progress_network = network.clone();
+ let result = match validate_blocks_extension(
+ local_ledger,
+ blocks,
+ adjusted_time_ms,
+ move |height| progress_network.update_sync_progress(progress_id, height),
+ )
+ .await
+ {
+ Ok(ledger) => network
+ .inner
+ .node
+ .lock()
+ .await
+ .import_verified_ledger(ledger)
+ .map(|_| ()),
+ Err(error) => Err(error),
+ };
+ drop(progress_guard);
let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error);
record_rejected_chain_payload(
network,
diff --git a/src/adapters/p2p/test_support.rs b/src/adapters/p2p/test_support.rs
@@ -55,6 +55,7 @@ pub(super) fn gossip_network(
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(std::sync::Mutex::new(InboundConnectionLimiter::default())),
metrics: P2pMetricsCounters::default(),
+ sync_progress: std::sync::Mutex::new(super::SyncProgressState::default()),
}),
}
}
diff --git a/src/adapters/p2p/tests.rs b/src/adapters/p2p/tests.rs
@@ -16,6 +16,35 @@ use tokio::io::AsyncWriteExt;
use super::test_support::{allocations, gossip_network, node, queue_plaintext_burn};
#[tokio::test]
+async fn block_batch_validation_reports_each_validated_height() {
+ let alice = Wallet::from_seed("batch-validation-progress-alice");
+ let allocations = allocations(std::slice::from_ref(&alice), 1_000);
+ let local = node("batch-progress-local", alice.clone(), allocations.clone());
+ let mut remote = node("batch-progress-remote", alice.clone(), allocations);
+ for timestamp_ms in [1, 2] {
+ queue_plaintext_burn(&mut remote, &alice, 1);
+ remote.drain_outbox();
+ remote.mine_one_at(timestamp_ms).unwrap();
+ remote.drain_outbox();
+ }
+ let blocks = remote.ledger().blocks_from(1, 10);
+ let progress = Arc::new(StdMutex::new(Vec::new()));
+ let reported = Arc::clone(&progress);
+
+ let adopted = super::validate_blocks_extension(
+ local.clone_ledger(),
+ blocks,
+ crate::app::now_ms(),
+ move |height| reported.lock().unwrap().push(height),
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(*progress.lock().unwrap(), vec![1, 2]);
+ assert_eq!(adopted.height(), 2);
+}
+
+#[tokio::test]
async fn full_outbound_queue_is_metric_not_peer_error() {
let wallet = Wallet::from_seed("full-outbound-queue");
let node = Arc::new(tokio::sync::Mutex::new(node(
@@ -37,6 +66,7 @@ async fn full_outbound_queue_is_metric_not_peer_error() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let (sender, _receiver) = tokio::sync::mpsc::channel(1);
@@ -137,9 +167,10 @@ async fn single_block_fork_error_requests_blocks_by_locator() {
let fork_blocks = remote_node.blocks_after_locator(&locator, limit);
assert_eq!(fork_blocks.len(), 2);
let local_ledger = network.inner.node.lock().await.clone_ledger();
- let adopted = super::validate_blocks_extension(local_ledger, fork_blocks, crate::app::now_ms())
- .await
- .unwrap();
+ let adopted =
+ super::validate_blocks_extension(local_ledger, fork_blocks, crate::app::now_ms(), |_| {})
+ .await
+ .unwrap();
assert_eq!(adopted.tip_hash(), remote_node.ledger().tip_hash());
}
@@ -371,6 +402,7 @@ async fn hello_rejects_wrong_network_or_genesis_without_banning() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
@@ -476,6 +508,7 @@ async fn hello_records_remote_clock_observation() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let remote_time_ms = crate::app::now_ms().saturating_add(60_000);
@@ -533,6 +566,7 @@ async fn hello_remembers_advertised_address_after_signed_session_and_dialback()
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let remote_node_id = super::new_node_id();
@@ -609,6 +643,7 @@ async fn hello_ignores_advertised_address_when_connected_peer_cannot_sign_claime
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let victim_node_id = super::new_node_id();
@@ -679,6 +714,7 @@ async fn dialback_rejects_address_that_signs_with_different_node_id() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let honest_node_id = super::new_node_id();
@@ -790,6 +826,7 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_bootstrap() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
@@ -871,6 +908,7 @@ async fn real_node_accepts_setup_placeholder_peer_and_pushes_bootstrap() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let setup_ledger = Ledger::new(BTreeMap::new(), 1);
@@ -946,6 +984,7 @@ async fn hello_ignores_private_advertised_listen_address() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let status = node.lock().await.ledger().status();
@@ -991,6 +1030,7 @@ async fn hello_ignores_loopback_alias_for_unspecified_self() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let hello = ProtocolHello {
@@ -1046,6 +1086,7 @@ async fn hello_removes_outbound_peer_that_announces_self_address() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let hello = ProtocolHello {
@@ -1100,6 +1141,7 @@ async fn hello_removes_outbound_peer_with_same_node_id() {
sessions: tokio::sync::Mutex::new(BTreeMap::new()),
inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())),
metrics: super::P2pMetricsCounters::default(),
+ sync_progress: StdMutex::new(super::SyncProgressState::default()),
}),
};
let hello = ProtocolHello {
diff --git a/src/main_tests.rs b/src/main_tests.rs
@@ -86,6 +86,8 @@ fn management_ui_blocks_interaction_while_the_node_is_syncing() {
assert!(html.contains("Synchronizing blockchain"));
assert!(html.contains("syncProgressPercent()"));
assert!(javascript.contains("this.networkHealth.state === \"syncing\""));
+ assert!(javascript.contains("this.networkHealth.sync_validated_height"));
+ assert!(javascript.contains("this.syncingNode() ? 1000 : 5000"));
assert!(javascript.contains("Syncing ${this.syncCurrentHeight().toLocaleString()} of"));
}
diff --git a/www/assets/iuna-ui.js b/www/assets/iuna-ui.js
@@ -158,9 +158,7 @@ window.iunaApp = function iunaApp() {
}
await this.refresh();
this.checkLatestRelease();
- if (!this.pollHandle) {
- this.pollHandle = setInterval(() => this.refresh({ silent: true }), 5000);
- }
+ this.schedulePoll();
},
canUseProtectedApi() {
@@ -169,10 +167,20 @@ window.iunaApp = function iunaApp() {
stopPolling() {
if (!this.pollHandle) return;
- clearInterval(this.pollHandle);
+ clearTimeout(this.pollHandle);
this.pollHandle = null;
},
+ schedulePoll() {
+ if (this.pollHandle || !this.canUseProtectedApi()) return;
+ const delay = this.syncingNode() ? 1000 : 5000;
+ this.pollHandle = setTimeout(async () => {
+ this.pollHandle = null;
+ await this.refresh({ silent: true });
+ this.schedulePoll();
+ }, delay);
+ },
+
tabFromHash() {
const hash = window.location.hash.replace(/^#\/?/, "");
return this.allowedTabs().includes(hash) ? hash : "wallet";
@@ -2976,12 +2984,23 @@ window.iunaApp = function iunaApp() {
},
syncCurrentHeight() {
- const height = Number(this.networkHealth.local_height ?? this.status.chain?.height ?? 0);
+ const height = Number(
+ this.networkHealth.sync_validated_height ??
+ this.networkHealth.local_height ??
+ this.status.chain?.height ??
+ 0
+ );
return Number.isFinite(height) && height >= 0 ? Math.floor(height) : 0;
},
syncTargetHeight() {
- const target = Number(this.networkHealth.best_known_height ?? this.syncCurrentHeight());
+ const batchTarget = Number(this.networkHealth.sync_target_height ?? 0);
+ const knownTarget = Number(this.networkHealth.best_known_height ?? 0);
+ const target = Math.max(
+ Number.isFinite(batchTarget) ? batchTarget : 0,
+ Number.isFinite(knownTarget) ? knownTarget : 0,
+ this.syncCurrentHeight()
+ );
return Number.isFinite(target) && target >= 0
? Math.max(this.syncCurrentHeight(), Math.floor(target))
: this.syncCurrentHeight();