Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
114 changes: 36 additions & 78 deletions crates/l1/src/subscriber.rs
Original file line number Diff line number Diff line change
Expand Up @@ -531,65 +531,31 @@ where
Ok(local_tempo_block_number + 1)
}

/// Resolve the first L1 block that has not already been ingested.
pub(crate) fn next_block_to_sync(&self) -> Result<u64, L1SubscriberError> {
let resolved = self.resolve_start_block()?;
let queued = self
.deposit_queue
.last_enqueued()
.map(|last| last.number.saturating_add(1));
let observed = self
.block_tracker
.latest()
.map(|last| last.number.saturating_add(1));

let next = [Some(resolved), queued, observed]
.into_iter()
.flatten()
.max()
.expect("resolved checkpoint is always present");

// Only the persisted zone checkpoint proves consumption.
// Queue and observation cursors are fetch high-water marks.
self.block_tracker
.initialize_consumed_through(resolved.saturating_sub(1));
Ok(next)
}

/// Return the block number referenced by the L1 `finalized` tag.
async fn finalized_block_number(
/// Synchronize all missing blocks through the current finalized L1 head.
///
/// The cursor advances after each block is fully applied.
pub(crate) async fn sync_finalized(
&self,
l1_provider: &impl Provider<TempoNetwork>,
) -> Result<u64, L1SubscriberError> {
Ok(l1_provider
next_block: &mut u64,
) -> Result<(), L1SubscriberError> {
let finalized = l1_provider
.get_header_by_number(BlockNumberOrTag::Finalized)
.await
.inspect_err(|_| self.subscriber_metrics.fetch_failures.increment(1))?
.map(|header| header.number())
.ok_or_eyre("L1 finalized block is not available")?)
}
.ok_or_eyre("L1 finalized block is not available")?;

/// Synchronize all missing blocks through the current finalized L1 head.
///
/// Callers provide the next block number and receive the next cursor after
/// a successful sync.
pub(crate) async fn sync_finalized_once(
&self,
l1_provider: &impl Provider<TempoNetwork>,
next_block: u64,
) -> Result<u64, L1SubscriberError> {
let finalized = self.finalized_block_number(l1_provider).await?;
if next_block > finalized {
self.record_seen_block(finalized, 0);
return Ok(next_block);
let pending_blocks = finalized.saturating_sub(next_block.saturating_sub(1));
self.record_seen_block(finalized, pending_blocks);
if pending_blocks == 0 {
return Ok(());
}

let blocks = finalized - next_block + 1;
self.record_seen_block(finalized, blocks);
info!(
from = next_block,
from = *next_block,
to = finalized,
blocks,
blocks = pending_blocks,
"Synchronizing finalized L1 blocks"
);

Expand All @@ -599,29 +565,7 @@ where
.backfill_duration_seconds
.record(start.elapsed().as_secs_f64());
self.subscriber_metrics.current_l1_lag_blocks.set(0.0);
Ok(finalized.saturating_add(1))
}

/// Follow finalized L1 using transport-specific head notifications as wakeups.
///
/// Header contents are intentionally ignored. Canonical block selection is
/// always based on the `finalized` tag read by [`Self::sync_finalized_once`].
pub(crate) async fn follow_finalized(
&self,
l1_provider: &impl Provider<TempoNetwork>,
mut stream: HeaderStream,
) -> Result<(), L1SubscriberError> {
let mut next_block = self.next_block_to_sync()?;

// Subscribe before the initial sync so a head published while catching
// up remains queued in the stream.
next_block = self.sync_finalized_once(l1_provider, next_block).await?;

while stream.next().await.is_some() {
next_block = self.sync_finalized_once(l1_provider, next_block).await?;
}

Err(eyre::eyre!("L1 head notification stream ended").into())
Ok(())
}

/// Backfill L1 blocks from `from..=to` with pipelined RPC fetching.
Expand All @@ -630,15 +574,16 @@ where
/// parallel, then processes them sequentially (event extraction and enqueue).
/// Receipts are fetched by the corresponding block
/// hash and validated against the header's receipts root before processing.
#[instrument(skip(self, l1_provider), fields(from, to))]
#[instrument(skip(self, l1_provider, next_block), fields(from = *next_block, to))]
async fn backfill(
&self,
l1_provider: &impl Provider<TempoNetwork>,
from: u64,
next_block: &mut u64,
to: u64,
) -> Result<(), L1SubscriberError> {
use futures::stream;

let from = *next_block;
let concurrency = self.config.l1_fetch_concurrency.max(1);
let subscriber_metrics = self.subscriber_metrics.clone();
let block_tracker = self.block_tracker.clone();
Expand Down Expand Up @@ -766,6 +711,7 @@ where
// configured retention sink and the contiguous observation tracker.
self.apply_enabled_token_events(&events);
self.update_l1_state_anchor(block_number, &invalidated);
*next_block = block_number.saturating_add(1);
if appended {
self.subscriber_metrics.blocks_enqueued.increment(1);
}
Expand Down Expand Up @@ -804,16 +750,28 @@ where
///
/// Transport and ordinary failures reconnect after the configured retry interval.
/// Deterministic failures while applying a receipt-verified finalized block are fatal.
pub async fn run(self) {
pub async fn run(self) -> Result<(), L1SubscriberError> {
let mut next_block = self.resolve_start_block()?;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ [ISSUE] Startup checkpoint read errors now bypass retry handling

run() now calls self.resolve_start_block()? before entering the reconnect/retry loop. Errors from provider.latest() or tempo_num_hash() are classified as retryable Other errors elsewhere, but here they are returned directly to the caller, which panics the critical task. A transient state-provider/storage read failure at subscriber startup can therefore crash the node instead of being retried as it was before this refactor.

Recommended Fix:
Initialize next_block through the same retry policy used for connection/sync errors, or add a small pre-loop retry around resolve_start_block() that retries provider/storage errors while still treating configuration errors such as an unanchored genesis as fatal.

self.block_tracker
.initialize_consumed_through(next_block.saturating_sub(1));

loop {
let result = async {
let result: Result<(), L1SubscriberError> = async {
let provider = self.connect().await?;
let header_stream = self.subscribe_block_headers(&provider).await?;
let mut header_stream = self.subscribe_block_headers(&provider).await?;
info!(
portal = %self.config.portal_address,
"Following finalized L1 blocks"
);
self.follow_finalized(&provider, header_stream).await

// Subscribe before the initial sync so a head published while catching
// up remains queued in the stream.
self.sync_finalized(&provider, &mut next_block).await?;
while let Some(()) = header_stream.next().await {
self.sync_finalized(&provider, &mut next_block).await?;
}

Err(eyre::eyre!("L1 head notification stream ended").into())
}
.await;

Expand All @@ -829,7 +787,7 @@ where
);
tokio::time::sleep(retry_interval).await;
} else {
panic!("{error}");
return Err(error);
}
}
}
Expand Down
83 changes: 65 additions & 18 deletions crates/l1/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -746,7 +746,7 @@ fn test_resolve_start_block_rejects_unanchored_genesis() {
}

#[tokio::test]
async fn test_follow_finalized_uses_new_heads_to_sync_missing_finalized_range() {
async fn test_sync_finalized_ingests_missing_finalized_range() {
let subscriber = test_subscriber(9);
let asserter = Asserter::new();
let l1_provider =
Expand All @@ -767,11 +767,16 @@ async fn test_follow_finalized_uses_new_heads_to_sync_missing_finalized_range()
push_header_and_empty_receipts(&asserter, header_11);
push_header_and_empty_receipts(&asserter, header_12);

let err = subscriber
.follow_finalized(&l1_provider, Box::pin(futures::stream::iter([()])))
let mut next_block = 10;
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap();
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.expect_err("finite header stream should end the subscriber");
assert!(err.to_string().contains("head notification stream ended"));
.unwrap();
assert_eq!(next_block, 13);

let blocks = subscriber.deposit_queue.drain();
assert_eq!(
Expand Down Expand Up @@ -813,23 +818,64 @@ async fn test_subscribe_block_headers_falls_back_to_http_block_filter() {
}

#[tokio::test]
async fn test_sync_finalized_once_does_not_refetch_current_cursor() {
async fn test_sync_finalized_does_not_refetch_current_cursor() {
let subscriber = test_subscriber(10);
let asserter = Asserter::new();
let l1_provider =
ProviderBuilder::new_with_network::<TempoNetwork>().connect_mocked_client(asserter.clone());
asserter.push_success(&Some(header_response(make_test_header(10))));

let next = subscriber
.sync_finalized_once(&l1_provider, 11)
let mut next_block = 11;
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap();

assert_eq!(next, 11);
assert_eq!(next_block, 11);
assert!(subscriber.deposit_queue.drain().is_empty());
assert!(asserter.read_q().is_empty());
}

#[tokio::test]
async fn test_sync_finalized_preserves_progress_after_partial_failure() {
let subscriber = test_subscriber(9);
let asserter = Asserter::new();
let l1_provider =
ProviderBuilder::new_with_network::<TempoNetwork>().connect_mocked_client(asserter.clone());
let header_10 = make_test_header(10);
let header_11 = make_chained_header(11, header_hash(&header_10));

asserter.push_success(&Some(header_response(header_11.clone())));
push_header_and_empty_receipts(&asserter, header_10);
asserter.push_failure_msg("temporary block fetch failure");

let mut next_block = 10;
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.expect_err("the first attempt should fail while fetching block 11");
assert_eq!(next_block, 11);

asserter.push_success(&Some(header_response(header_11.clone())));
push_header_and_empty_receipts(&asserter, header_11);
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap();

assert_eq!(next_block, 12);
assert_eq!(
subscriber
.deposit_queue
.drain()
.iter()
.map(|block| block.header.number())
.collect::<Vec<_>>(),
vec![10, 11]
);
assert!(asserter.read_q().is_empty());
}

#[test]
fn test_push_log_decodes_withdrawal_bounce_back() {
let portal_address = address!("0x0000000000000000000000000000000000000ABC");
Expand Down Expand Up @@ -1730,8 +1776,9 @@ async fn sync_classifies_corrupt_recognized_portal_log_as_fatal() {
asserter.push_success(&Some(header_response(header_10)));
asserter.push_success(&Some(vec![receipt]));

let mut next_block = 10;
let err = subscriber
.sync_finalized_once(&l1_provider, 10)
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap_err();
assert!(matches!(
Expand Down Expand Up @@ -1797,13 +1844,12 @@ async fn sync_applies_leadership_transition_before_enqueueing_the_activation_blo
asserter.push_success(&Some(header_response(header_10.clone())));
asserter.push_success(&Some(vec![receipt]));

assert_eq!(
subscriber
.sync_finalized_once(&l1_provider, 10)
.await
.unwrap(),
11
);
let mut next_block = 10;
subscriber
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap();
assert_eq!(next_block, 11);

let seen = sink.seen.lock();
assert_eq!(seen.len(), 1);
Expand Down Expand Up @@ -1843,8 +1889,9 @@ async fn sync_fails_fatally_when_the_leadership_sink_rejects_the_transition() {
asserter.push_success(&Some(header_response(header_10.clone())));
asserter.push_success(&Some(vec![receipt]));

let mut next_block = 10;
let err = subscriber
.sync_finalized_once(&l1_provider, 10)
.sync_finalized(&l1_provider, &mut next_block)
.await
.unwrap_err();
assert!(matches!(
Expand Down
10 changes: 9 additions & 1 deletion crates/node/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -662,7 +662,15 @@ where
self.encryption_keys.clone(),
);
let task_executor = ctx.node.task_executor().clone();
task_executor.spawn_critical_task("l1-block-subscriber", Box::pin(l1_subscriber.run()));
task_executor.spawn_critical_task(
"l1-block-subscriber",
Box::pin(async move {
l1_subscriber
.run()
.await
.unwrap_or_else(|error| panic!("{error}"));
}),
);
info!(target: "reth::cli", "L1 subscriber started with deposit enqueueing");

// Start the Commonware network and the long-lived event router
Expand Down