Skip to content

Commit 29ab04d

Browse files
committed
wip
1 parent dc4bcab commit 29ab04d

3 files changed

Lines changed: 124 additions & 35 deletions

File tree

crates/consensus/src/follow/executor/actor.rs

Lines changed: 62 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -183,15 +183,20 @@ where
183183
}
184184

185185
let request = if let Some((block, ack)) = self.block_queue.pop_front() {
186-
let finalized_advanced = self
187-
.forkchoice
188-
.advance_finalized(Target::from_block(&block));
189-
let forkchoice = (self.forkchoice.requires_update() || heartbeat)
186+
let block_target = Target::from_block(&block);
187+
let finality_anchor =
188+
self.forkchoice
189+
.advance_finalized(block_target)
190+
.then_some(ForkchoiceTargets {
191+
head: block_target,
192+
finalized: block_target,
193+
});
194+
let latest_forkchoice = (self.forkchoice.requires_update() || heartbeat)
190195
.then_some(self.forkchoice.latest());
191196
ExecutionRequest::Block {
192197
block,
193-
forkchoice,
194-
finalized_advanced,
198+
latest_forkchoice,
199+
finality_anchor,
195200
ack,
196201
}
197202
} else if self.forkchoice.requires_update() || heartbeat {
@@ -258,12 +263,14 @@ where
258263
enum ExecutionRequest {
259264
Forkchoice(ForkchoiceTargets),
260265
Block {
266+
/// A finalized block delivered by marshal for execution.
261267
block: Block,
262-
forkchoice: Option<ForkchoiceTargets>,
268+
/// The latest combined head and finalized targets, if an update is due.
269+
latest_forkchoice: Option<ForkchoiceTargets>,
270+
/// The delivered block as both head and finalized when it advances finality.
271+
finality_anchor: Option<ForkchoiceTargets>,
272+
/// Signals marshal after the payload and required forkchoice update complete.
263273
ack: Exact,
264-
/// Whether the block advances finality and must wait for a `VALID` FCU before
265-
/// acknowledgement.
266-
finalized_advanced: bool,
267274
},
268275
}
269276

@@ -294,34 +301,64 @@ async fn execute_request<TContext: Pacer, E: ExecutionEngine + 'static>(
294301
}
295302
ExecutionRequest::Block {
296303
block,
297-
forkchoice,
298-
finalized_advanced,
304+
latest_forkchoice,
305+
finality_anchor,
299306
ack,
300307
} => {
301308
if let Err(error) = submit_new_payload(&context, &execution_engine, block).await {
302309
return ExecutionTaskResult::Fatal(error);
303310
}
304311

305-
if let Some(forkchoice) = forkchoice {
306-
loop {
312+
let submitted_fcu = match latest_forkchoice {
313+
Some(latest_forkchoice) => {
307314
let result =
308-
submit_forkchoice_update(&context, &execution_engine, &forkchoice).await;
315+
submit_forkchoice_update(&context, &execution_engine, &latest_forkchoice)
316+
.await;
309317
match result {
310-
Ok(ForkchoiceOutcome::Valid) => break,
311-
Ok(ForkchoiceOutcome::Syncing) if finalized_advanced => {
312-
debug!(
313-
"execution layer is syncing before applying finalized FCU; retrying"
314-
);
315-
context.sleep(FINALIZED_FCU_RETRY_INTERVAL).await;
316-
}
317-
Ok(ForkchoiceOutcome::Syncing) => break,
318+
Ok(ForkchoiceOutcome::Valid) => Some(latest_forkchoice),
319+
Ok(ForkchoiceOutcome::Syncing) => match finality_anchor {
320+
Some(finality_anchor) => {
321+
if let Err(error) = submit_finality_anchor(
322+
&context,
323+
&execution_engine,
324+
&finality_anchor,
325+
)
326+
.await
327+
{
328+
return ExecutionTaskResult::Fatal(error);
329+
}
330+
Some(finality_anchor)
331+
}
332+
None => Some(latest_forkchoice),
333+
},
318334
Err(error) => return ExecutionTaskResult::Fatal(error),
319335
}
320336
}
321-
}
337+
None => None,
338+
};
322339

323340
ack.acknowledge();
324-
ExecutionTaskResult::Completed(forkchoice)
341+
ExecutionTaskResult::Completed(submitted_fcu)
342+
}
343+
}
344+
}
345+
346+
async fn submit_finality_anchor<TContext: Pacer, E: ExecutionEngine + ?Sized>(
347+
context: &TContext,
348+
execution_engine: &E,
349+
finality_anchor: &ForkchoiceTargets,
350+
) -> eyre::Result<()> {
351+
loop {
352+
let outcome = submit_forkchoice_update(context, execution_engine, finality_anchor).await?;
353+
match outcome {
354+
ForkchoiceOutcome::Valid => return Ok(()),
355+
ForkchoiceOutcome::Syncing => {
356+
debug!(
357+
"execution layer is syncing before applying finalized block; retrying block \
358+
anchor FCU"
359+
);
360+
context.sleep(FINALIZED_FCU_RETRY_INTERVAL).await;
361+
}
325362
}
326363
}
327364
}

crates/consensus/src/follow/executor/test/mod.rs

Lines changed: 55 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -421,7 +421,7 @@ fn finalization_waits_for_durable_block_before_advancing_finalized() {
421421
}
422422

423423
#[test_traced]
424-
fn syncing_block_forkchoice_retries_before_acknowledging() {
424+
fn syncing_certificate_head_falls_back_to_block_anchor_before_acknowledging() {
425425
deterministic::Runner::default().start(|context| async move {
426426
let provider = StubExecutionProvider::default();
427427
let release_head_forkchoice = provider.pause_next_forkchoice();
@@ -443,6 +443,7 @@ fn syncing_block_forkchoice_retries_before_acknowledging() {
443443
mailbox.finalization(round(2), future_head);
444444
wait_until(&context, || provider.forkchoices().len() == 1).await;
445445

446+
provider.set_syncing_forkchoice_head(future_head.0);
446447
provider.set_forkchoices_syncing(true);
447448
let block = make_block_at_round(1, B256::ZERO, round(1));
448449
let finalized = block.block_hash();
@@ -452,25 +453,70 @@ fn syncing_block_forkchoice_retries_before_acknowledging() {
452453
release_head_forkchoice
453454
.send(())
454455
.expect("the head FCU should still be waiting");
455-
wait_until(&context, || provider.forkchoices().len() == 2).await;
456+
wait_until(&context, || provider.forkchoices().len() == 3).await;
456457

457458
let mut waiter = Box::pin(waiter);
458459
tokio::select! {
459-
result = &mut waiter => panic!("syncing FCU acknowledged the block: {result:?}"),
460+
result = &mut waiter => panic!("syncing block anchor acknowledged the block: {result:?}"),
460461
_ = context.sleep(Duration::from_millis(100)) => {}
461462
}
462463

463464
provider.set_forkchoices_syncing(false);
464465
waiter
465466
.await
466-
.expect("valid retry should acknowledge the durable block");
467+
.expect("valid block anchor should acknowledge the durable block");
468+
469+
wait_until(&context, || provider.forkchoices().len() == 5).await;
467470

468471
let forkchoices = provider.forkchoices();
469-
assert_eq!(forkchoices.len(), 3);
470-
assert_eq!(forkchoices[1], forkchoices[2]);
471-
assert_eq!(forkchoices[2].head_block_hash, future_head.0);
472-
assert_eq!(forkchoices[2].safe_block_hash, finalized);
473-
assert_eq!(forkchoices[2].finalized_block_hash, finalized);
472+
assert_eq!(forkchoices[1].head_block_hash, future_head.0);
473+
assert_eq!(forkchoices[1].safe_block_hash, finalized);
474+
assert_eq!(forkchoices[1].finalized_block_hash, finalized);
475+
476+
assert_eq!(forkchoices[2], forkchoices[3]);
477+
assert_eq!(forkchoices[3].head_block_hash, finalized);
478+
assert_eq!(forkchoices[3].safe_block_hash, finalized);
479+
assert_eq!(forkchoices[3].finalized_block_hash, finalized);
480+
481+
assert_eq!(forkchoices[4], forkchoices[1]);
482+
});
483+
}
484+
485+
#[test_traced]
486+
fn known_certificate_head_finalizes_delivered_block_without_anchor() {
487+
deterministic::Runner::default().start(|context| async move {
488+
let provider = StubExecutionProvider::default();
489+
490+
let (actor, mut mailbox) = init(
491+
context.child("follower_executor"),
492+
Config {
493+
execution_provider: provider.clone(),
494+
execution_engine: provider.clone(),
495+
marshal: StubMarshal::default(),
496+
epoch_strategy: FixedEpocher::new(EPOCH_LENGTH),
497+
floor: Height::zero(),
498+
fcu_heartbeat_interval: Duration::from_secs(60),
499+
},
500+
);
501+
actor.start();
502+
503+
let future_head = digest(9);
504+
mailbox.finalization(round(2), future_head);
505+
wait_until(&context, || provider.forkchoices().len() == 1).await;
506+
507+
let block = make_block_at_round(1, B256::ZERO, round(1));
508+
let finalized = block.block_hash();
509+
let (ack, waiter) = Exact::handle();
510+
assert!(mailbox.report(Update::Block(block.into(), ack)).accepted());
511+
waiter
512+
.await
513+
.expect("known certificate head should finalize the delivered block");
514+
515+
let forkchoices = provider.forkchoices();
516+
assert_eq!(forkchoices.len(), 2);
517+
assert_eq!(forkchoices[1].head_block_hash, future_head.0);
518+
assert_eq!(forkchoices[1].safe_block_hash, finalized);
519+
assert_eq!(forkchoices[1].finalized_block_hash, finalized);
474520
});
475521
}
476522

crates/consensus/src/follow/executor/test/utils.rs

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ struct StubExecutionProviderInner {
7171
reject_payloads: AtomicBool,
7272
reject_forkchoices: AtomicBool,
7373
sync_forkchoices: AtomicBool,
74+
syncing_forkchoice_head: Mutex<Option<B256>>,
7475
forkchoice_gate: Mutex<Option<oneshot::Receiver<()>>>,
7576
}
7677

@@ -107,6 +108,10 @@ impl StubExecutionProvider {
107108
self.inner.sync_forkchoices.store(syncing, Ordering::SeqCst);
108109
}
109110

111+
pub(super) fn set_syncing_forkchoice_head(&self, head: B256) {
112+
*self.inner.syncing_forkchoice_head.lock() = Some(head);
113+
}
114+
110115
pub(super) fn pause_next_forkchoice(&self) -> oneshot::Sender<()> {
111116
let (release, gate) = oneshot::channel();
112117
*self.inner.forkchoice_gate.lock() = Some(gate);
@@ -183,7 +188,8 @@ impl ExecutionEngine for StubExecutionProvider {
183188
self.inner.forkchoices.lock().push(state);
184189
let gate = self.inner.forkchoice_gate.lock().take();
185190
let rejected = self.inner.reject_forkchoices.load(Ordering::SeqCst);
186-
let syncing = self.inner.sync_forkchoices.load(Ordering::SeqCst);
191+
let syncing = self.inner.sync_forkchoices.load(Ordering::SeqCst)
192+
|| *self.inner.syncing_forkchoice_head.lock() == Some(state.head_block_hash);
187193
async move {
188194
if let Some(gate) = gate {
189195
let _ = gate.await;

0 commit comments

Comments
 (0)