commit 1ae1f0259835f81f8de4185354c54ab925fe28a1
parent ffeb7105be4062c7e54e331d2d41349d0d5e9833
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Sun, 23 Aug 2026 22:22:14 +0200
Cancel stale VDF work
Diffstat:
8 files changed, 347 insertions(+), 65 deletions(-)
diff --git a/src/app/receive.rs b/src/app/receive.rs
@@ -1,7 +1,8 @@
use anyhow::Result;
use crate::domain::{
- Block, BurnBundle, ChainSnapshot, Ledger, Transaction, TransactionSubmitOutcome,
+ Block, BurnBundle, ChainSnapshot, Ledger, MINE_ANCHOR_LIMIT_REACHED, Transaction,
+ TransactionSubmitOutcome,
};
use super::{GossipEnvelope, IMPORT_REBROADCAST_LIMIT, NodeCore};
@@ -13,11 +14,17 @@ impl NodeCore {
}
pub fn receive_gossiped_transaction(&mut self, tx: Transaction) -> Result<()> {
- if self
- .ledger
- .submit_transaction_with_outcome(tx.clone())?
- .added()
- {
+ let outcome = match self.ledger.submit_transaction_with_outcome(tx.clone()) {
+ Ok(outcome) => outcome,
+ Err(error)
+ if matches!(&tx, Transaction::Mine { .. })
+ && error.to_string() == MINE_ANCHOR_LIMIT_REACHED =>
+ {
+ return Ok(());
+ }
+ Err(error) => return Err(error),
+ };
+ if outcome.added() {
self.outbox.push(GossipEnvelope::Transaction(tx));
}
Ok(())
@@ -211,6 +218,41 @@ mod tests {
}
#[test]
+ fn gossiped_mine_over_anchor_limit_is_silently_ignored() {
+ let wallets = (0..4)
+ .map(|index| {
+ let seed = format!("gossip-mine-limit-{index}");
+ Wallet::from_seed(&seed)
+ })
+ .collect::<Vec<_>>();
+ let ledger = Ledger::new(BTreeMap::new(), 1);
+ let mine_actions = wallets
+ .iter()
+ .map(|wallet| ledger.build_mine(wallet.address()).unwrap())
+ .collect::<Vec<_>>();
+ let mut node = NodeCore::from_ledger(wallets[0].clone(), ledger, 0);
+
+ node.receive_gossiped_transaction(mine_actions[0].clone())
+ .unwrap();
+ node.receive_gossiped_transaction(mine_actions[1].clone())
+ .unwrap();
+ node.receive_gossiped_transaction(mine_actions[2].clone())
+ .unwrap();
+
+ assert_eq!(node.ledger().pending().len(), 2);
+ assert!(
+ node.ledger()
+ .pending()
+ .iter()
+ .all(|transaction| transaction.signature() != mine_actions[2].signature())
+ );
+ let error = node
+ .receive_transaction(mine_actions[3].clone())
+ .unwrap_err();
+ assert_eq!(error.to_string(), "mine transaction anchor limit reached");
+ }
+
+ #[test]
fn burn_bundle_imports_new_signed_burns_to_mempool() {
let alice = Wallet::from_seed("bundle-import-new-burn-alice");
let bob = Wallet::from_seed("bundle-import-new-burn-bob");
diff --git a/src/domain.rs b/src/domain.rs
@@ -47,6 +47,7 @@ use ledger_ops::{
apply_transaction, credit_reward_output, ensure_block_has_burn, ensure_block_has_burn_from,
recovery_vdf_seed_for_child, validate_genesis_burn_transaction, vdf_seed_for_child,
};
+pub(crate) use ledger_pending::MINE_ANCHOR_LIMIT_REACHED;
pub use ledger_state::Ledger;
use ledger_state::unix_now_ms;
use mining::{mine_payload, mine_signature};
@@ -80,7 +81,10 @@ pub use validation::validate_address;
use validation::{
canonical_transaction_size_bytes, validate_hash, validate_protocol_id, validate_signature,
};
-pub use vdf::{VdfProgress, VdfProgressPhase, run_vdf, run_vdf_with_progress, verify_vdf};
+pub use vdf::{
+ VdfProgress, VdfProgressPhase, run_vdf, run_vdf_cancellable_with_progress,
+ run_vdf_with_progress, verify_vdf,
+};
pub use wallet::Wallet;
pub fn burn_committee_slot_count(eligible_rank_count: usize) -> usize {
diff --git a/src/domain/block.rs b/src/domain/block.rs
@@ -233,6 +233,10 @@ impl PreparedBlock {
self.height
}
+ pub fn prev_hash(&self) -> &str {
+ &self.prev_hash
+ }
+
pub fn timestamp_ms(&self) -> u64 {
self.timestamp_ms
}
diff --git a/src/domain/ledger_pending.rs b/src/domain/ledger_pending.rs
@@ -24,6 +24,8 @@ use super::{
MINE_ACTIONS_PER_ANCHOR_LIMIT, OutPoint, Transaction, TxOutput,
};
+pub(crate) const MINE_ANCHOR_LIMIT_REACHED: &str = "mine transaction anchor limit reached";
+
impl Ledger {
pub(super) fn valid_pending_transactions(&self) -> Vec<Transaction> {
let mut utxos = self.utxos.clone();
@@ -345,7 +347,7 @@ impl Ledger {
.count(),
);
if known_count >= MINE_ACTIONS_PER_ANCHOR_LIMIT {
- bail!("mine transaction anchor limit reached");
+ bail!(MINE_ANCHOR_LIMIT_REACHED);
}
}
Ok(())
diff --git a/src/domain/vdf/mod.rs b/src/domain/vdf/mod.rs
@@ -1,4 +1,7 @@
-use std::time::{Duration, Instant};
+use std::{
+ sync::atomic::AtomicBool,
+ time::{Duration, Instant},
+};
use super::{Block, FinalizerMode, MAX_VDF_ROUNDS, VDF_TARGET_BLOCK_MS, decode_hex, hex_encode};
@@ -45,30 +48,47 @@ pub fn run_vdf_with_progress(
seed: &str,
rounds: u64,
progress_interval: Duration,
- mut progress: impl FnMut(VdfProgress),
+ progress: impl FnMut(VdfProgress),
) -> String {
+ let cancelled = AtomicBool::new(false);
+ run_vdf_cancellable_with_progress(seed, rounds, progress_interval, &cancelled, progress)
+ .expect("non-cancellable VDF must finish")
+}
+
+pub fn run_vdf_cancellable_with_progress(
+ seed: &str,
+ rounds: u64,
+ progress_interval: Duration,
+ cancelled: &AtomicBool,
+ mut progress: impl FnMut(VdfProgress),
+) -> Option<String> {
let total_steps = rounds.saturating_mul(2);
let mut last_progress = Instant::now();
- let solution = wesolowski::prove(seed.as_bytes(), rounds, |phase, completed_phase_rounds| {
- let completed_steps = match phase {
- VdfProgressPhase::Output => completed_phase_rounds,
- VdfProgressPhase::Proof => rounds.saturating_add(completed_phase_rounds),
- };
- maybe_report_vdf_progress(
- &mut last_progress,
- progress_interval,
- VdfProgress {
- completed_steps,
- total_steps,
- completed_phase_rounds,
- phase_rounds: rounds,
- phase,
- },
- &mut progress,
- );
- })
+ let solution = wesolowski::prove_cancellable(
+ seed.as_bytes(),
+ rounds,
+ |phase, completed_phase_rounds| {
+ let completed_steps = match phase {
+ VdfProgressPhase::Output => completed_phase_rounds,
+ VdfProgressPhase::Proof => rounds.saturating_add(completed_phase_rounds),
+ };
+ maybe_report_vdf_progress(
+ &mut last_progress,
+ progress_interval,
+ VdfProgress {
+ completed_steps,
+ total_steps,
+ completed_phase_rounds,
+ phase_rounds: rounds,
+ phase,
+ },
+ &mut progress,
+ );
+ },
+ cancelled,
+ )
.expect("valid IUNA VDF parameters must produce a class-group proof");
- encode_vdf_solution(&solution)
+ solution.map(|solution| encode_vdf_solution(&solution))
}
pub fn verify_vdf(seed: &str, rounds: u64, solution: &str) -> bool {
@@ -151,11 +171,14 @@ fn decode_vdf_solution(solution: &str) -> Option<Vec<u8>> {
#[cfg(test)]
mod tests {
- use std::time::Duration;
+ use std::{
+ sync::atomic::{AtomicBool, Ordering},
+ time::Duration,
+ };
use super::{
- VDF_SOLUTION_PREFIX, VdfProgressPhase, run_vdf, run_vdf_with_progress,
- vdf_solution_placeholder, verify_vdf, wesolowski,
+ VDF_SOLUTION_PREFIX, VdfProgressPhase, run_vdf, run_vdf_cancellable_with_progress,
+ run_vdf_with_progress, vdf_solution_placeholder, verify_vdf, wesolowski,
};
#[test]
@@ -197,6 +220,26 @@ mod tests {
}
#[test]
+ fn cancellable_vdf_stops_during_output_or_proof() {
+ for phase in [VdfProgressPhase::Output, VdfProgressPhase::Proof] {
+ let cancelled = AtomicBool::new(false);
+ let solution = run_vdf_cancellable_with_progress(
+ "cancelled-seed",
+ 128,
+ Duration::ZERO,
+ &cancelled,
+ |progress| {
+ if progress.phase == phase && progress.completed_phase_rounds >= 10 {
+ cancelled.store(true, Ordering::Relaxed);
+ }
+ },
+ );
+
+ assert!(solution.is_none(), "VDF did not stop during {phase:?}");
+ }
+ }
+
+ #[test]
fn vdf_solution_uses_chia_bqfc_protocol_format() {
let solution = run_vdf("test-seed", 16);
let encoded = solution.strip_prefix(VDF_SOLUTION_PREFIX).unwrap();
diff --git a/src/domain/vdf/prover.rs b/src/domain/vdf/prover.rs
@@ -1,3 +1,5 @@
+use std::sync::atomic::{AtomicBool, Ordering};
+
use kyn_vdf::{Form, KynVdfError, get_b};
use num_bigint::{BigInt, BigUint};
use num_traits::{One, ToPrimitive};
@@ -24,6 +26,15 @@ struct ClassGroup<'a> {
threshold: &'a BigInt,
}
+struct CheckpointProofInput<'a> {
+ generator: &'a Form,
+ output: &'a Form,
+ checkpoints: &'a [limb_arithmetic::LimbForm],
+ rounds: u64,
+ parameters: ProofParameters,
+}
+
+#[cfg(test)]
pub(super) fn prove(
discriminant: &BigInt,
generator: &Form,
@@ -31,25 +42,53 @@ pub(super) fn prove(
rounds: u64,
progress: impl FnMut(VdfProgressPhase, u64),
) -> Result<(Form, Form), KynVdfError> {
+ let cancelled = AtomicBool::new(false);
+ prove_cancellable(
+ discriminant,
+ generator,
+ threshold,
+ rounds,
+ progress,
+ &cancelled,
+ )?
+ .ok_or_else(|| arithmetic_error("non-cancellable VDF was cancelled"))
+}
+
+pub(super) fn prove_cancellable(
+ discriminant: &BigInt,
+ generator: &Form,
+ threshold: &BigInt,
+ rounds: u64,
+ progress: impl FnMut(VdfProgressPhase, u64),
+ cancelled: &AtomicBool,
+) -> Result<Option<(Form, Form)>, KynVdfError> {
let group = ClassGroup {
discriminant,
threshold,
};
let parameters = ProofParameters::for_rounds(rounds);
if parameters.checkpoint_count > MAX_CHECKPOINTS || parameters.bucket_count > MAX_BUCKETS {
- return prove_constant_memory(discriminant, generator, threshold, rounds, progress);
+ return prove_constant_memory_cancellable(
+ discriminant,
+ generator,
+ threshold,
+ rounds,
+ progress,
+ cancelled,
+ );
}
- prove_checkpointed(group, generator, rounds, parameters, progress)
+ prove_checkpointed_cancellable(group, generator, rounds, parameters, progress, cancelled)
}
-fn prove_checkpointed(
+fn prove_checkpointed_cancellable(
group: ClassGroup<'_>,
generator: &Form,
rounds: u64,
parameters: ProofParameters,
mut progress: impl FnMut(VdfProgressPhase, u64),
-) -> Result<(Form, Form), KynVdfError> {
+ cancelled: &AtomicBool,
+) -> Result<Option<(Form, Form)>, KynVdfError> {
let checkpoint_capacity = usize::try_from(parameters.checkpoint_count)
.map_err(|_| arithmetic_error("checkpoint count does not fit in memory"))?;
let checkpoint_stride = u64::from(parameters.k)
@@ -61,6 +100,9 @@ fn prove_checkpointed(
let mut output = limb_arithmetic::LimbForm::from_form(generator);
let mut output_scratch = limb_arithmetic::LimbFormScratch::default();
for completed_rounds in 1..=rounds {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
if (completed_rounds - 1) % checkpoint_stride == 0 {
checkpoints.push(output.clone());
}
@@ -74,27 +116,37 @@ fn prove_checkpointed(
let output = output.into_form();
debug_assert_eq!(checkpoints.len(), checkpoint_capacity);
- let proof = generate_checkpoint_proof(
+ let Some(proof) = generate_checkpoint_proof_cancellable(
group,
- generator,
- &output,
- &checkpoints,
- rounds,
- parameters,
+ CheckpointProofInput {
+ generator,
+ output: &output,
+ checkpoints: &checkpoints,
+ rounds,
+ parameters,
+ },
&mut progress,
- )?;
- Ok((output, proof))
+ cancelled,
+ )?
+ else {
+ return Ok(None);
+ };
+ Ok(Some((output, proof)))
}
-fn generate_checkpoint_proof(
+fn generate_checkpoint_proof_cancellable(
group: ClassGroup<'_>,
- generator: &Form,
- output: &Form,
- checkpoints: &[limb_arithmetic::LimbForm],
- rounds: u64,
- parameters: ProofParameters,
+ input: CheckpointProofInput<'_>,
progress: &mut impl FnMut(VdfProgressPhase, u64),
-) -> Result<Form, KynVdfError> {
+ cancelled: &AtomicBool,
+) -> Result<Option<Form>, KynVdfError> {
+ let CheckpointProofInput {
+ generator,
+ output,
+ checkpoints,
+ rounds,
+ parameters,
+ } = input;
let challenge = get_b(group.discriminant, generator, output)?;
let bucket_count = usize::try_from(parameters.bucket_count)
.map_err(|_| arithmetic_error("bucket count does not fit in memory"))?;
@@ -122,6 +174,9 @@ fn generate_checkpoint_proof(
let block_step = BigUint::from(2_u8).modpow(&BigUint::from(block_step_exponent), &challenge);
for j in (0..parameters.l).rev() {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
proof = proof.fast_pow_u64_with_scratch(
1_u64 << parameters.k,
&limb_discriminant,
@@ -133,6 +188,9 @@ fn generate_checkpoint_proof(
get_blocks_for_pass(j, parameters, rounds, &challenge, &block_step)?;
let mut buckets: Vec<Option<limb_arithmetic::LimbForm>> = vec![None; bucket_count];
for (i, checkpoint) in checkpoints.iter().enumerate() {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
let bucket = checkpoint_blocks[i];
if bucket != INVALID_BUCKET {
let bucket_form = buckets[bucket].take();
@@ -152,6 +210,9 @@ fn generate_checkpoint_proof(
}
for b1 in 0..row_count {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
let row_start = b1 << k0;
let mut aggregate: Option<limb_arithmetic::LimbForm> = None;
for b0 in 0..column_count {
@@ -188,6 +249,9 @@ fn generate_checkpoint_proof(
}
for b0 in 0..column_count {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
let mut aggregate: Option<limb_arithmetic::LimbForm> = None;
for b1 in 0..row_count {
if let Some(bucket) = &buckets[((b1 << k0) + b0) as usize] {
@@ -225,7 +289,7 @@ fn generate_checkpoint_proof(
proof.reduce();
progress(VdfProgressPhase::Proof, rounds);
- Ok(proof.into_form())
+ Ok(Some(proof.into_form()))
}
fn get_blocks_for_pass(
@@ -311,18 +375,22 @@ fn block_from_residue(
.ok_or_else(|| arithmetic_error("proof block does not fit in its bucket range"))
}
-fn prove_constant_memory(
+fn prove_constant_memory_cancellable(
discriminant: &BigInt,
generator: &Form,
threshold: &BigInt,
rounds: u64,
mut progress: impl FnMut(VdfProgressPhase, u64),
-) -> Result<(Form, Form), KynVdfError> {
+ cancelled: &AtomicBool,
+) -> Result<Option<(Form, Form)>, KynVdfError> {
let limb_discriminant = limb_arithmetic::to_limb(discriminant);
let limb_threshold = limb_arithmetic::to_limb(threshold);
let mut output = limb_arithmetic::LimbForm::from_form(generator);
let mut output_scratch = limb_arithmetic::LimbFormScratch::default();
for completed_rounds in 1..=rounds {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
output = output.nudupl_reduce_with_scratch(
&limb_discriminant,
&limb_threshold,
@@ -338,6 +406,9 @@ fn prove_constant_memory(
let mut proof_scratch = limb_arithmetic::LimbFormScratch::default();
let mut remainder = BigUint::one() % &challenge;
for completed_rounds in 1..=rounds {
+ if cancelled.load(Ordering::Relaxed) {
+ return Ok(None);
+ }
let doubled = &remainder << 1_usize;
let carry = doubled >= challenge;
proof = proof.nudupl_reduce_with_scratch(
@@ -357,7 +428,66 @@ fn prove_constant_memory(
progress(VdfProgressPhase::Proof, completed_rounds);
}
- Ok((output, proof.into_form()))
+ Ok(Some((output, proof.into_form())))
+}
+
+#[cfg(test)]
+fn prove_checkpointed(
+ group: ClassGroup<'_>,
+ generator: &Form,
+ rounds: u64,
+ parameters: ProofParameters,
+ progress: impl FnMut(VdfProgressPhase, u64),
+) -> Result<(Form, Form), KynVdfError> {
+ let cancelled = AtomicBool::new(false);
+ prove_checkpointed_cancellable(group, generator, rounds, parameters, progress, &cancelled)?
+ .ok_or_else(|| arithmetic_error("non-cancellable VDF was cancelled"))
+}
+
+#[cfg(test)]
+fn generate_checkpoint_proof(
+ group: ClassGroup<'_>,
+ generator: &Form,
+ output: &Form,
+ checkpoints: &[limb_arithmetic::LimbForm],
+ rounds: u64,
+ parameters: ProofParameters,
+ progress: &mut impl FnMut(VdfProgressPhase, u64),
+) -> Result<Form, KynVdfError> {
+ let cancelled = AtomicBool::new(false);
+ generate_checkpoint_proof_cancellable(
+ group,
+ CheckpointProofInput {
+ generator,
+ output,
+ checkpoints,
+ rounds,
+ parameters,
+ },
+ progress,
+ &cancelled,
+ )?
+ .ok_or_else(|| arithmetic_error("non-cancellable VDF was cancelled"))
+}
+
+#[cfg(test)]
+fn prove_constant_memory(
+ discriminant: &BigInt,
+ generator: &Form,
+ threshold: &BigInt,
+ rounds: u64,
+ progress: impl FnMut(VdfProgressPhase, u64),
+) -> Result<(Form, Form), KynVdfError> {
+ let cancelled = AtomicBool::new(false);
+ prove_constant_memory_cancellable(
+ discriminant,
+ generator,
+ threshold,
+ rounds,
+ progress,
+ &cancelled,
+ )?
+ .ok_or_else(|| arithmetic_error("non-cancellable VDF was cancelled"))
}
fn arithmetic_error(message: &str) -> KynVdfError {
diff --git a/src/domain/vdf/wesolowski.rs b/src/domain/vdf/wesolowski.rs
@@ -1,3 +1,5 @@
+use std::sync::atomic::AtomicBool;
+
use kyn_vdf::{
Form, KynVdfError, create_discriminant, deserialize_form, isqrt_fourth, serialize_form,
verify_wesolowski,
@@ -10,11 +12,24 @@ pub(super) const DISCRIMINANT_BITS: usize = 1024;
const FORM_BYTES: usize = 100;
pub(super) const SOLUTION_BYTES: usize = FORM_BYTES * 2;
+#[cfg(test)]
pub(super) fn prove(
seed: &[u8],
rounds: u64,
- mut progress: impl FnMut(VdfProgressPhase, u64),
+ progress: impl FnMut(VdfProgressPhase, u64),
) -> Result<Vec<u8>, KynVdfError> {
+ let cancelled = AtomicBool::new(false);
+ prove_cancellable(seed, rounds, progress, &cancelled)?.ok_or_else(|| {
+ KynVdfError::ArithmeticError("non-cancellable VDF was cancelled".to_string())
+ })
+}
+
+pub(super) fn prove_cancellable(
+ seed: &[u8],
+ rounds: u64,
+ mut progress: impl FnMut(VdfProgressPhase, u64),
+ cancelled: &AtomicBool,
+) -> Result<Option<Vec<u8>>, KynVdfError> {
if rounds == 0 {
return Err(KynVdfError::InvalidIterations(rounds));
}
@@ -24,9 +39,18 @@ pub(super) fn prove(
Form::generator(&discriminant).ok_or(KynVdfError::InvalidDiscriminantIdentity)?;
let threshold = isqrt_fourth(&discriminant.abs());
- let (output, proof) =
- prover::prove(&discriminant, &generator, &threshold, rounds, &mut progress)?;
- serialize_solution(&output, &proof)
+ let Some((output, proof)) = prover::prove_cancellable(
+ &discriminant,
+ &generator,
+ &threshold,
+ rounds,
+ &mut progress,
+ cancelled,
+ )?
+ else {
+ return Ok(None);
+ };
+ serialize_solution(&output, &proof).map(Some)
}
pub(super) fn verify(seed: &[u8], rounds: u64, solution: &[u8]) -> bool {
diff --git a/src/main.rs b/src/main.rs
@@ -2,7 +2,10 @@ use std::{
collections::BTreeMap,
net::SocketAddr,
path::{Path, PathBuf},
- sync::Arc,
+ sync::{
+ Arc,
+ atomic::{AtomicBool, Ordering},
+ },
time::{Duration, Instant},
};
@@ -18,7 +21,8 @@ use iuna::{
},
domain::{
Amount, ChainSnapshot, GenesisBurn, LaunchProfile, Ledger, MAX_VDF_ROUNDS, MICRO_IUNA,
- VDF_TARGET_BLOCK_MS, VdfProgress, VdfProgressPhase, run_vdf, run_vdf_with_progress,
+ VDF_TARGET_BLOCK_MS, VdfProgress, VdfProgressPhase, run_vdf,
+ run_vdf_cancellable_with_progress,
},
};
use tokio::sync::Mutex;
@@ -656,6 +660,7 @@ async fn run_automatic_finalizer(node: SharedNode, gossip: p2p::GossipNetwork, d
}
let candidate_height = work.height();
+ let candidate_parent = work.prev_hash().to_string();
let seed = work.vdf_seed().to_string();
let rounds = work.vdf_rounds();
let publish_at_ms = work.timestamp_ms();
@@ -675,11 +680,20 @@ async fn run_automatic_finalizer(node: SharedNode, gossip: p2p::GossipNetwork, d
continue;
}
let (progress_tx, progress_rx) = std::sync::mpsc::channel();
+ let cancellation = Arc::new(AtomicBool::new(false));
+ let worker_cancellation = Arc::clone(&cancellation);
let mut vdf_worker = tokio::task::spawn_blocking(move || {
- run_vdf_with_progress(&seed, rounds, VDF_PROGRESS_LOG_INTERVAL, |progress| {
- let _ = progress_tx.send(progress);
- })
+ run_vdf_cancellable_with_progress(
+ &seed,
+ rounds,
+ VDF_PROGRESS_LOG_INTERVAL,
+ worker_cancellation.as_ref(),
+ |progress| {
+ let _ = progress_tx.send(progress);
+ },
+ )
});
+ let mut cancelled_for_new_tip = false;
let vdf_output = loop {
tokio::select! {
result = &mut vdf_worker => {
@@ -689,7 +703,7 @@ async fn run_automatic_finalizer(node: SharedNode, gossip: p2p::GossipNetwork, d
if debug {
eprintln!("VDF worker failed: {error:#}");
}
- continue;
+ None
}
};
}
@@ -701,9 +715,28 @@ async fn run_automatic_finalizer(node: SharedNode, gossip: p2p::GossipNetwork, d
}
node.lock().await.record_automatic_finalization_status(message);
}
+ let tip_changed = node.lock().await.ledger().tip_hash() != candidate_parent;
+ if tip_changed {
+ cancelled_for_new_tip = true;
+ cancellation.store(true, Ordering::Relaxed);
+ }
}
}
};
+ let Some(vdf_output) = vdf_output else {
+ let message = if cancelled_for_new_tip {
+ format!("cancelled stale VDF for candidate block {candidate_height}")
+ } else {
+ format!("VDF worker failed for candidate block {candidate_height}")
+ };
+ if debug {
+ println!("auto-finalization {message}");
+ }
+ node.lock()
+ .await
+ .record_automatic_finalization_status(message);
+ continue;
+ };
let completed_at_ms = now_ms();
let publish_timestamp_ms = completed_at_ms.max(publish_at_ms);