Skip to content
Closed
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
2 changes: 1 addition & 1 deletion crates/l1/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ pub use event::{EnabledToken, L1PortalEvents, L1SequencerEvent};
pub use ext::{ChainTempoStateExt, TempoStateExt};
pub use queue::DepositQueue;
pub use state::{L1StateCache, PolicyCache, PolicyProvider};
pub use subscriber::{L1Subscriber, L1SubscriberConfig};
pub use subscriber::{L1BlockObserver, L1Subscriber, L1SubscriberConfig};

pub(crate) use event::EnqueueOutcome;

Expand Down
289 changes: 233 additions & 56 deletions crates/l1/src/subscriber.rs

Large diffs are not rendered by default.

115 changes: 114 additions & 1 deletion crates/l1/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,11 +149,12 @@ fn test_subscriber(
genesis_tempo_block_number,
policy_cache: crate::PolicyCache::default(),
l1_state_cache: crate::L1StateCache::new(HashSet::from([portal_address])),
block_observer: L1BlockObserver::default(),
l1_fetch_concurrency: 1,
retry_connection_interval: Duration::from_secs(1),
},
local_state,
deposit_queue: DepositQueue::default(),
deposit_queue: Some(DepositQueue::default()),
tracked_tokens: vec![],
tip403_metrics: Default::default(),
subscriber_metrics: Default::default(),
Expand Down Expand Up @@ -451,6 +452,118 @@ fn update_l1_state_anchor_reorg_clears_stale_policy_state() {
);
}

#[tokio::test]
async fn l1_block_observer_waits_for_exact_hash_and_invalidates_descendants() {
let observer = L1BlockObserver::default();
let hash_10 = B256::with_last_byte(0x10);
let hash_11 = B256::with_last_byte(0x11);
observer.record(NumHash::new(10, hash_10));
observer.record(NumHash::new(11, hash_11));

observer.wait_for(NumHash::new(11, hash_11)).await.unwrap();
assert!(
observer
.wait_for(NumHash::new(11, B256::with_last_byte(0xff)))
.await
.unwrap_err()
.to_string()
.contains("observed different L1 hash")
);

let replacement_10 = B256::with_last_byte(0xa0);
observer.record(NumHash::new(10, replacement_10));
assert_eq!(observer.observed_hash(10), Some(replacement_10));
assert_eq!(observer.observed_hash(11), None);
}

#[tokio::test]
async fn l1_block_observer_prune_drops_consumed_anchors() {
let observer = L1BlockObserver::default();
for number in 10..=14 {
observer.record(NumHash::new(number, B256::with_last_byte(number as u8)));
}

observer.prune_below(12);
assert_eq!(observer.observed_hash(11), None);
assert_eq!(observer.observed_hash(12), Some(B256::with_last_byte(12)));
assert_eq!(
observer.latest(),
Some(NumHash::new(14, B256::with_last_byte(14)))
);

// Waiting on a pruned height fails instead of hanging.
assert!(
observer
.wait_for(NumHash::new(11, B256::with_last_byte(11)))
.await
.unwrap_err()
.to_string()
.contains("below the observer's retained range")
);
}

#[tokio::test]
async fn observed_policy_change_is_visible_at_its_exact_l1_height() {
use crate::state::tip403::AuthRole;
use tempo_contracts::precompiles::ITIP403Registry::PolicyType;

let subscriber = test_subscriber(
Arc::new(SequenceLocalTempoCheckpointReader::new([0])),
Some(0),
);
let token = address!("0x0000000000000000000000000000000000000011");
let user = address!("0x0000000000000000000000000000000000000022");
{
let mut cache = subscriber.config.policy_cache.write();
cache.set_policy_type(2, PolicyType::WHITELIST);
cache.set_token_policy(token, 10, 2);
}

subscriber.apply_policy_events(
10,
&[PolicyEvent::MembershipChanged {
policy_id: 2,
account: user,
in_set: true,
}],
);
let hash_10 = B256::with_last_byte(0x10);
subscriber
.config
.block_observer
.record(NumHash::new(10, hash_10));

subscriber.apply_policy_events(
11,
&[PolicyEvent::MembershipChanged {
policy_id: 2,
account: user,
in_set: false,
}],
);
let hash_11 = B256::with_last_byte(0x11);
subscriber
.config
.block_observer
.record(NumHash::new(11, hash_11));
subscriber
.config
.block_observer
.wait_for(NumHash::new(11, hash_11))
.await
.unwrap();

let cache = subscriber.config.policy_cache.read();
assert_eq!(
cache.is_authorized(token, user, 10, AuthRole::Transfer),
Some(true)
);
assert_eq!(
cache.is_authorized(token, user, 11, AuthRole::Transfer),
Some(false)
);
}

/// Confirm the front of the queue, panicking if it fails.
fn confirm(queue: &mut PendingDeposits) -> L1BlockDeposits {
let num_hash = queue.peek().expect("queue is empty").header.num_hash();
Expand Down
56 changes: 36 additions & 20 deletions crates/node/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ use tracing::{debug, info, warn};
use zone_chainspec::ZoneChainSpec;
use zone_evm::ZoneEvmConfig;
use zone_l1::{
DepositQueue, L1Subscriber, L1SubscriberConfig, PolicyCache, TempoStateExt,
DepositQueue, L1BlockObserver, L1Subscriber, L1SubscriberConfig, PolicyCache, TempoStateExt,
state::{
L1StateCache, L1StateProvider, L1StateProviderConfig, PolicyProvider,
spawn_policy_resolution_task, spawn_pool_prefetch_task,
Expand Down Expand Up @@ -214,6 +214,7 @@ impl ZoneNode {
genesis_tempo_block_number,
policy_cache: policy_cache.clone(),
l1_state_cache: l1_state_cache.clone(),
block_observer: L1BlockObserver::default(),
l1_fetch_concurrency,
retry_connection_interval,
};
Expand Down Expand Up @@ -309,6 +310,11 @@ impl ZoneNode {
self.policy_cache.clone()
}

/// Returns the shared record of independently observed L1 anchors.
pub fn l1_block_observer(&self) -> L1BlockObserver {
self.l1_config.block_observer.clone()
}

/// Returns a [`ComponentsBuilder`] configured for a Zone node.
pub fn components<N>(
executor_builder: ZoneExecutorBuilder,
Expand Down Expand Up @@ -467,17 +473,7 @@ where
.erased();

self.resolve_and_seed_tokens(&l1_provider).await?;
let p2p_role = self.p2p_config.as_ref().map(P2pConfig::role);
if p2p_role == Some(Role::Follower) {
// TODO(multi-sequencer): Split L1 observation/cache updates from deposit
// enqueueing. Followers import complete blocks from the leader and do not consume
// DepositQueue; starting the unified subscriber here would grow that queue forever.
// On promotion/restart the subscriber resumes from the tempoBlockNumber persisted in
// the follower's imported zone state.
info!(target: "reth::cli", "Skipping L1 deposit subscriber on follower");
} else {
self.spawn_l1_subscriber(&ctx);
}
self.spawn_l1_subscriber(&ctx);
self.spawn_policy_tasks(&l1_provider, &ctx);

let task_executor = ctx.node.task_executor().clone();
Expand All @@ -490,6 +486,8 @@ where
&task_executor,
ctx.node.provider().clone(),
ctx.beacon_engine_handle.clone(),
self.l1_config.block_observer.clone(),
self.policy_cache.clone(),
)?;
}

Expand Down Expand Up @@ -546,6 +544,8 @@ where
task_executor: &reth_tasks::TaskExecutor,
provider: N::Provider,
engine: reth_node_builder::ConsensusEngineHandle<ZonePayloadTypes>,
l1_observer: L1BlockObserver,
policy_cache: PolicyCache,
) -> eyre::Result<()> {
let role = config.role();
let handle = spawn_p2p(config, network_id)?;
Expand All @@ -572,7 +572,14 @@ where
// for later ACK/backfill commands even though followers send nothing in this PR.
task_executor.spawn_critical_task(
"zone-p2p-block-import",
import_leader_blocks(provider, engine, events, commands),
import_leader_blocks(
provider,
engine,
events,
commands,
l1_observer,
policy_cache,
),
);
}
}
Expand Down Expand Up @@ -671,13 +678,22 @@ where

/// Spawn the L1 subscriber. Listens for new blocks and deposit events.
fn spawn_l1_subscriber(&mut self, ctx: &AddOnsContext<'_, N>) {
L1Subscriber::spawn(
self.l1_config.clone(),
ctx.node.provider().clone(),
self.deposit_queue.clone(),
ctx.node.task_executor().clone(),
);
info!(target: "reth::cli", "Unified L1 subscriber started");
if self.p2p_config.as_ref().map(P2pConfig::role) == Some(Role::Follower) {
L1Subscriber::spawn_observer(
self.l1_config.clone(),
ctx.node.provider().clone(),
ctx.node.task_executor().clone(),
);
info!(target: "reth::cli", "L1 observer started without deposit enqueueing");
} else {
L1Subscriber::spawn(
self.l1_config.clone(),
ctx.node.provider().clone(),
self.deposit_queue.clone(),
ctx.node.task_executor().clone(),
);
info!(target: "reth::cli", "Leader L1 observer and deposit sink started");
}
}

/// Spawn TIP-403 policy resolution and pool prefetch tasks.
Expand Down
Loading
Loading