diff --git a/crates/supervisor/core/src/chain_processor/task.rs b/crates/supervisor/core/src/chain_processor/task.rs
index 792eca941e..a2209cdb31 100644
--- a/crates/supervisor/core/src/chain_processor/task.rs
+++ b/crates/supervisor/core/src/chain_processor/task.rs
@@ -18,7 +18,7 @@ use tracing::{debug, error, info};
/// It listens for events emitted by the managed node and handles them accordingly.
#[derive(Debug)]
pub struct ChainProcessorTask
{
- _rollup_config: RollupConfig,
+ rollup_config: RollupConfig,
chain_id: ChainId,
metrics_enabled: Option,
@@ -54,7 +54,7 @@ where
) -> Self {
let log_indexer = LogIndexer::new(managed_node.clone(), state_manager.clone());
Self {
- _rollup_config: rollup_config,
+ rollup_config,
chain_id,
metrics_enabled: None,
cancel_token,
@@ -326,37 +326,53 @@ where
block_number = derived_ref_pair.derived.number,
"Processing local safe derived block pair"
);
- match self.state_manager.save_derived_block(derived_ref_pair) {
- Ok(_) => Ok(derived_ref_pair.derived),
- Err(StorageError::BlockOutOfOrder | StorageError::ConflictError(_)) => {
- error!(
- target: "chain_processor",
- chain_id = self.chain_id,
- block_number = derived_ref_pair.derived.number,
- "Block out of order detected, resetting managed node"
- );
- if let Err(err) = self.managed_node.reset().await {
+ if self.rollup_config.is_post_interop(derived_ref_pair.derived.timestamp) {
+ match self.state_manager.save_derived_block(derived_ref_pair) {
+ Ok(_) => return Ok(derived_ref_pair.derived),
+ Err(StorageError::BlockOutOfOrder) => {
error!(
target: "chain_processor",
chain_id = self.chain_id,
+ block_number = derived_ref_pair.derived.number,
+ "Block out of order detected, resetting managed node"
+ );
+
+ if let Err(err) = self.managed_node.reset().await {
+ error!(
+ target: "chain_processor",
+ chain_id = self.chain_id,
+ %err,
+ "Failed to reset managed node after block out of order"
+ );
+ }
+ return Err(StorageError::BlockOutOfOrder.into());
+ }
+ Err(err) => {
+ error!(
+ target: "chain_processor",
+ chain_id = self.chain_id,
+ block_number = derived_ref_pair.derived.number,
%err,
- "Failed to reset managed node after block out of order"
+ "Failed to save derived block pair"
);
+ return Err(err.into());
}
- Err(StorageError::BlockOutOfOrder.into())
- }
- Err(err) => {
- error!(
- target: "chain_processor",
- chain_id = self.chain_id,
- block_number = derived_ref_pair.derived.number,
- %err,
- "Failed to save derived block pair"
- );
- Err(err.into())
}
}
+
+ if self.rollup_config.is_interop_activation_block(derived_ref_pair.derived) {
+ info!(
+ target: "chain_processor",
+ chain_id = self.chain_id,
+ block_number = derived_ref_pair.derived.number,
+ "Initialising derivation storage for interop activation block"
+ );
+ self.state_manager.initialise_derivation_storage(derived_ref_pair)?;
+ return Ok(derived_ref_pair.derived);
+ }
+
+ Ok(derived_ref_pair.derived)
}
async fn handle_unsafe_event(
@@ -370,7 +386,21 @@ where
"Processing unsafe block"
);
- self.log_indexer.clone().sync_logs(block);
+ if self.rollup_config.is_post_interop(block.timestamp) {
+ self.log_indexer.clone().sync_logs(block);
+ return Ok(block);
+ }
+
+ if self.rollup_config.is_interop_activation_block(block) {
+ info!(
+ target: "chain_processor",
+ chain_id = self.chain_id,
+ block_number = block.number,
+ "Initialising log storage for interop activation block"
+ );
+ self.state_manager.initialise_log_storage(block)?;
+ return Ok(block);
+ }
Ok(block)
}
@@ -412,6 +442,7 @@ where
mod tests {
use super::*;
use crate::{
+ config::Genesis,
event::ChainEvent,
syncnode::{
BlockProvider, ManagedNodeController, ManagedNodeDataProvider, ManagedNodeError,
@@ -548,13 +579,60 @@ mod tests {
}
);
+ fn genesis() -> Genesis {
+ let l2 = BlockInfo::new(B256::from([1u8; 32]), 0, B256::ZERO, 50);
+ let l1 = BlockInfo::new(B256::from([2u8; 32]), 10, B256::ZERO, 1000);
+ Genesis::new(l1, l2)
+ }
+
+ fn get_rollup_config(interop_time: u64) -> RollupConfig {
+ RollupConfig::new(genesis(), 2, Some(interop_time))
+ }
+
#[tokio::test]
- async fn test_handle_unsafe_event_triggers() {
+ async fn test_handle_unsafe_event_pre_interop() {
+ let mockdb = MockDb::new();
+ let mocknode = MockNode::new();
+
+ // Send unsafe block event
+ let block = BlockInfo::new(B256::ZERO, 123, B256::ZERO, 10);
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let rollup_config = get_rollup_config(1000);
+
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ tx.send(ChainEvent::UnsafeBlock { block }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+
+ // Give it time to process
+ tokio::time::sleep(Duration::from_millis(50)).await;
+
+ // Stop the task
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_handle_unsafe_event_post_interop() {
let mut mockdb = MockDb::new();
let mut mocknode = MockNode::new();
// Send unsafe block event
- let block = BlockInfo::new(B256::ZERO, 123, B256::ZERO, 0);
+ let block = BlockInfo::new(B256::ZERO, 123, B256::ZERO, 1003);
mockdb.expect_store_block_logs().returning(move |_block, _log| Ok(()));
mocknode.expect_fetch_receipts().returning(move |block_hash| {
@@ -568,7 +646,8 @@ mod tests {
let cancel_token = CancellationToken::new();
let (tx, rx) = mpsc::channel(10);
- let rollup_config = RollupConfig::default();
+ let rollup_config = get_rollup_config(1000);
+
let task = ChainProcessorTask::new(
rollup_config,
1,
@@ -591,7 +670,45 @@ mod tests {
}
#[tokio::test]
- async fn test_handle_derived_event_triggers() {
+ async fn test_handle_unsafe_event_interop_activation() {
+ let mut mockdb = MockDb::new();
+ let mocknode = MockNode::new();
+
+ // Block that triggers interop activation
+ let block = BlockInfo::new(B256::ZERO, 123, B256::ZERO, 1001); // Use timestamp/number that triggers activation
+
+ let rollup_config = get_rollup_config(1000);
+
+ mockdb.expect_initialise_log_storage().returning(move |b| {
+ assert_eq!(b, block);
+ Ok(())
+ });
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ tx.send(ChainEvent::UnsafeBlock { block }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+ tokio::time::sleep(std::time::Duration::from_millis(50)).await;
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_handle_derived_event_pre_interop() {
let block_pair = DerivedRefPair {
source: BlockInfo {
number: 123,
@@ -603,8 +720,57 @@ mod tests {
number: 1234,
hash: B256::ZERO,
parent_hash: B256::ZERO,
+ timestamp: 999,
+ },
+ };
+
+ let mockdb = MockDb::new();
+ let mocknode = MockNode::new();
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let rollup_config = get_rollup_config(1000);
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ // Send unsafe block event
+ tx.send(ChainEvent::DerivedBlock { derived_ref_pair: block_pair }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+
+ // Give it time to process
+ tokio::time::sleep(Duration::from_millis(50)).await;
+
+ // Stop the task
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_handle_derived_event_post_interop() {
+ let block_pair = DerivedRefPair {
+ source: BlockInfo {
+ number: 123,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
timestamp: 0,
},
+ derived: BlockInfo {
+ number: 1234,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 1003,
+ },
};
let mut mockdb = MockDb::new();
@@ -621,7 +787,7 @@ mod tests {
let cancel_token = CancellationToken::new();
let (tx, rx) = mpsc::channel(10);
- let rollup_config = RollupConfig::default();
+ let rollup_config = get_rollup_config(1000);
let task = ChainProcessorTask::new(
rollup_config,
1,
@@ -644,6 +810,164 @@ mod tests {
task_handle.await.unwrap();
}
+ #[tokio::test]
+ async fn test_handle_derived_event_interop_activation() {
+ let block_pair = DerivedRefPair {
+ source: BlockInfo {
+ number: 123,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 0,
+ },
+ derived: BlockInfo {
+ number: 1234,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 1001,
+ },
+ };
+
+ let mut mockdb = MockDb::new();
+ let mocknode = MockNode::new();
+
+ mockdb.expect_initialise_derivation_storage().returning(move |_pair: DerivedRefPair| {
+ assert_eq!(_pair, block_pair);
+ Ok(())
+ });
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let rollup_config = get_rollup_config(1000);
+
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ // Send unsafe block event
+ tx.send(ChainEvent::DerivedBlock { derived_ref_pair: block_pair }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+
+ // Give it time to process
+ tokio::time::sleep(Duration::from_millis(50)).await;
+
+ // Stop the task
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_handle_derived_event_block_out_of_order_triggers_reset() {
+ let block_pair = DerivedRefPair {
+ source: BlockInfo {
+ number: 123,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 0,
+ },
+ derived: BlockInfo {
+ number: 1234,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 1003, // post-interop
+ },
+ };
+
+ let mut mockdb = MockDb::new();
+ let mut mocknode = MockNode::new();
+
+ // Simulate BlockOutOfOrder error
+ mockdb
+ .expect_save_derived_block()
+ .returning(move |_pair: DerivedRefPair| Err(StorageError::BlockOutOfOrder));
+
+ // Expect reset to be called
+ mocknode.expect_reset().returning(|| Ok(()));
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let rollup_config = get_rollup_config(1000);
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ tx.send(ChainEvent::DerivedBlock { derived_ref_pair: block_pair }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+
+ tokio::time::sleep(std::time::Duration::from_millis(50)).await;
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_handle_derived_event_other_error() {
+ let block_pair = DerivedRefPair {
+ source: BlockInfo {
+ number: 123,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 0,
+ },
+ derived: BlockInfo {
+ number: 1234,
+ hash: B256::ZERO,
+ parent_hash: B256::ZERO,
+ timestamp: 1003, // post-interop
+ },
+ };
+
+ let mut mockdb = MockDb::new();
+ let mocknode = MockNode::new();
+
+ // Simulate a different error
+ mockdb
+ .expect_save_derived_block()
+ .returning(move |_pair: DerivedRefPair| Err(StorageError::DatabaseNotInitialised));
+
+ let writer = Arc::new(mockdb);
+ let managed_node = Arc::new(mocknode);
+
+ let cancel_token = CancellationToken::new();
+ let (tx, rx) = mpsc::channel(10);
+
+ let rollup_config = get_rollup_config(1000);
+ let task = ChainProcessorTask::new(
+ rollup_config,
+ 1,
+ managed_node,
+ writer,
+ cancel_token.clone(),
+ rx,
+ );
+
+ tx.send(ChainEvent::DerivedBlock { derived_ref_pair: block_pair }).await.unwrap();
+
+ let task_handle = tokio::spawn(task.run());
+
+ tokio::time::sleep(std::time::Duration::from_millis(50)).await;
+ cancel_token.cancel();
+ task_handle.await.unwrap();
+ }
+
#[tokio::test]
async fn test_handle_derivation_origin_update_triggers() {
let origin =
diff --git a/crates/supervisor/core/src/config/rollup_config_set.rs b/crates/supervisor/core/src/config/rollup_config_set.rs
index 3ce6437484..8f6dfd6da5 100644
--- a/crates/supervisor/core/src/config/rollup_config_set.rs
+++ b/crates/supervisor/core/src/config/rollup_config_set.rs
@@ -29,8 +29,8 @@ impl Genesis {
}
}
- /// Returns the genesis anchor as a [`DerivedRefPair`].
- pub const fn get_anchor(&self) -> DerivedRefPair {
+ /// Returns the genesis as a [`DerivedRefPair`].
+ pub const fn get_derived_pair(&self) -> DerivedRefPair {
DerivedRefPair { derived: self.l2, source: self.l1 }
}
}
@@ -73,15 +73,36 @@ impl RollupConfig {
})
}
+ /// Returns `true` if the timestamp is at or after the interop activation time.
+ ///
+ /// Interop activates at [`interop_time`](Self::interop_time). This function checks whether the
+ /// provided timestamp is before or after interop timestamp.
+ ///
+ /// Returns `false` if `interop_time` is not configured.
+ pub fn is_interop(&self, timestamp: u64) -> bool {
+ self.interop_time.is_some_and(|t| timestamp >= t)
+ }
+
/// Returns `true` if the timestamp is strictly after the interop activation block.
///
/// Interop activates at [`interop_time`](Self::interop_time). This function checks whether the
- /// current block timestamp is *after* that activation, skipping the activation block
+ /// provided timestamp is *after* that activation, skipping the activation block
/// itself.
///
/// Returns `false` if `interop_time` is not configured.
pub fn is_post_interop(&self, timestamp: u64) -> bool {
- self.interop_time.is_some_and(|t| timestamp.saturating_sub(self.block_time) >= t)
+ self.is_interop(timestamp.saturating_sub(self.block_time))
+ }
+
+ /// Returns `true` if given block is the interop activation block.
+ ///
+ /// An interop activation block is defined as the block that is right after the
+ /// interop activation time.
+ ///
+ /// Returns `false` if `interop_time` is not configured.
+ pub fn is_interop_activation_block(&self, block: BlockInfo) -> bool {
+ self.is_interop(block.timestamp) &&
+ !self.is_interop(block.timestamp.saturating_sub(self.block_time))
}
}
@@ -115,10 +136,15 @@ impl RollupConfigSet {
Ok(())
}
- /// returns whether interop is enabled for a chain at given timestamp
+ /// Returns `true` if interop is enabled for the chain at given timestamp.
pub fn is_interop_enabled(&self, chain_id: ChainId, timestamp: u64) -> bool {
self.get(chain_id).map(|cfg| cfg.is_post_interop(timestamp)).unwrap_or(false) // if config not found, return false
}
+
+ /// Returns `true` if given block is the interop activation block for the specified chain.
+ pub fn is_interop_activation_block(&self, chain_id: ChainId, block: BlockInfo) -> bool {
+ self.get(chain_id).map(|cfg| cfg.is_interop_activation_block(block)).unwrap_or(false)
+ }
}
#[cfg(test)]
@@ -152,4 +178,14 @@ mod tests {
// Unknown chain_id returns false
assert!(!set.is_interop_enabled(ChainId::from(999u64), 200));
}
+
+ #[test]
+ fn test_rollup_config_is_interop_interop_time_zero() {
+ // Interop time is 100, block_time is 10
+ let rollup_config =
+ RollupConfig::new(Genesis::new(dummy_blockinfo(0), dummy_blockinfo(0)), 2, Some(0));
+
+ assert!(rollup_config.is_interop(0));
+ assert!(rollup_config.is_interop(1000));
+ }
}
diff --git a/crates/supervisor/core/src/logindexer/indexer.rs b/crates/supervisor/core/src/logindexer/indexer.rs
index 02eac64c06..42302e676c 100644
--- a/crates/supervisor/core/src/logindexer/indexer.rs
+++ b/crates/supervisor/core/src/logindexer/indexer.rs
@@ -272,7 +272,7 @@ mod tests {
let mut mock_provider = MockBlockProvider::new();
mock_provider.expect_fetch_receipts().withf(move |hash| *hash == block_hash).returning(
|_| {
- Err(ManagedNodeError::Client(ClientError::Authentication(
+ Err(ManagedNodeError::ClientError(ClientError::Authentication(
AuthenticationError::InvalidHeader,
)))
},
diff --git a/crates/supervisor/core/src/safety_checker/task.rs b/crates/supervisor/core/src/safety_checker/task.rs
index adae55f6e9..5446130220 100644
--- a/crates/supervisor/core/src/safety_checker/task.rs
+++ b/crates/supervisor/core/src/safety_checker/task.rs
@@ -6,7 +6,7 @@ use crate::{
use alloy_primitives::ChainId;
use derive_more::Constructor;
use kona_protocol::BlockInfo;
-use kona_supervisor_storage::CrossChainSafetyProvider;
+use kona_supervisor_storage::{CrossChainSafetyProvider, StorageError};
use std::{sync::Arc, time::Duration};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
@@ -114,10 +114,27 @@ where
// Finds the next block that is eligible for promotion at the configured target level.
fn find_next_promotable_block(&self) -> Result {
- let current_head =
- self.provider.get_safety_head_ref(self.chain_id, self.promoter.target_level())?;
- let upper_head =
- self.provider.get_safety_head_ref(self.chain_id, self.promoter.lower_bound_level())?;
+ let current_head = self
+ .provider
+ .get_safety_head_ref(self.chain_id, self.promoter.target_level())
+ .map_err(|err| {
+ if matches!(err, StorageError::FutureData) {
+ CrossSafetyError::NoBlockToPromote
+ } else {
+ err.into()
+ }
+ })?;
+
+ let upper_head = self
+ .provider
+ .get_safety_head_ref(self.chain_id, self.promoter.lower_bound_level())
+ .map_err(|err| {
+ if matches!(err, StorageError::FutureData) {
+ CrossSafetyError::NoBlockToPromote
+ } else {
+ err.into()
+ }
+ })?;
if current_head.number >= upper_head.number {
return Err(CrossSafetyError::NoBlockToPromote);
diff --git a/crates/supervisor/core/src/supervisor.rs b/crates/supervisor/core/src/supervisor.rs
index 25653e15d5..cfa5ddede3 100644
--- a/crates/supervisor/core/src/supervisor.rs
+++ b/crates/supervisor/core/src/supervisor.rs
@@ -143,9 +143,13 @@ impl Supervisor {
for (chain_id, config) in self.config.rollup_config_set.rollups.iter() {
// Initialise the database for each chain.
let db = self.database_factory.get_or_create_db(*chain_id)?;
- let anchor = config.genesis.get_anchor();
- db.initialise_log_storage(anchor.derived)?;
- db.initialise_derivation_storage(anchor)?;
+ let interop_time = config.interop_time;
+ let derived_pair = config.genesis.get_derived_pair();
+ if config.is_interop(derived_pair.derived.timestamp) {
+ info!(target: "supervisor_service", chain_id, interop_time, %derived_pair, "Initialising database for interop activation block");
+ db.initialise_log_storage(derived_pair.derived)?;
+ db.initialise_derivation_storage(derived_pair)?;
+ }
info!(target: "supervisor_service", chain_id, "Database initialized successfully");
}
Ok(())
diff --git a/crates/supervisor/core/src/syncnode/client.rs b/crates/supervisor/core/src/syncnode/client.rs
index 8706c944a3..9fc0b731f7 100644
--- a/crates/supervisor/core/src/syncnode/client.rs
+++ b/crates/supervisor/core/src/syncnode/client.rs
@@ -42,6 +42,9 @@ pub trait ManagedNodeClient: Debug {
/// Fetches the [`BlockInfo`] by block number.
async fn block_ref_by_number(&self, block_number: u64) -> Result;
+ /// Resets the managed node to the pre-interop state.
+ async fn reset_pre_interop(&self) -> Result<(), ClientError>;
+
/// Resets the node state with the provided block IDs.
async fn reset(
&self,
@@ -302,7 +305,7 @@ impl ManagedNodeClient for Client {
Metrics::MANAGED_NODE_RPC_REQUEST_DURATION_SECONDS,
"block_ref_by_number",
async {
- ManagedModeApiClient::block_ref_by_number(client.as_ref(), block_number).await
+ ManagedModeApiClient::l2_block_ref_by_number(client.as_ref(), block_number).await
},
"node" => self.config.url.clone()
)?;
@@ -310,6 +313,21 @@ impl ManagedNodeClient for Client {
Ok(block_info)
}
+ async fn reset_pre_interop(&self) -> Result<(), ClientError> {
+ let client = self.get_ws_client().await?;
+ observe_metrics_for_result_async!(
+ Metrics::MANAGED_NODE_RPC_REQUESTS_SUCCESS_TOTAL,
+ Metrics::MANAGED_NODE_RPC_REQUESTS_ERROR_TOTAL,
+ Metrics::MANAGED_NODE_RPC_REQUEST_DURATION_SECONDS,
+ "reset_pre_interop",
+ async {
+ ManagedModeApiClient::reset_pre_interop(client.as_ref()).await
+ },
+ "node" => self.config.url.clone()
+ )?;
+ Ok(())
+ }
+
async fn reset(
&self,
unsafe_id: BlockNumHash,
diff --git a/crates/supervisor/core/src/syncnode/error.rs b/crates/supervisor/core/src/syncnode/error.rs
index 6aab59f503..01d1233f45 100644
--- a/crates/supervisor/core/src/syncnode/error.rs
+++ b/crates/supervisor/core/src/syncnode/error.rs
@@ -7,7 +7,7 @@ use thiserror::Error;
pub enum ManagedNodeError {
/// Represents an error that occurred while starting the managed node.
#[error(transparent)]
- Client(#[from] ClientError),
+ ClientError(#[from] ClientError),
/// Represents an error that occurred while subscribing to the managed node.
#[error("subscription error: {0}")]
diff --git a/crates/supervisor/core/src/syncnode/resetter.rs b/crates/supervisor/core/src/syncnode/resetter.rs
index f7f0646b50..4c27c9b810 100644
--- a/crates/supervisor/core/src/syncnode/resetter.rs
+++ b/crates/supervisor/core/src/syncnode/resetter.rs
@@ -1,6 +1,6 @@
use super::{ManagedNodeClient, ManagedNodeError};
use alloy_eips::BlockNumHash;
-use kona_supervisor_storage::{DerivationStorageReader, HeadRefStorageReader};
+use kona_supervisor_storage::{DerivationStorageReader, HeadRefStorageReader, StorageError};
use kona_supervisor_types::SuperHead;
use std::sync::Arc;
use tokio::sync::Mutex;
@@ -32,6 +32,11 @@ where
let SuperHead { local_unsafe, cross_unsafe, local_safe, cross_safe, finalized, .. } =
match self.get_latest_valid_super_head().await {
Ok(block) => block,
+ // todo: require refactor and corner case handling
+ Err(ManagedNodeError::StorageError(StorageError::DatabaseNotInitialised)) => {
+ self.reset_pre_interop().await?;
+ return Ok(());
+ }
Err(err) => {
error!(target: "resetter", %err, "Failed to get latest valid derived block");
return Err(ManagedNodeError::ResetFailed);
@@ -59,7 +64,15 @@ where
.inspect_err(|err| {
error!(target: "resetter", %err, "Failed to reset managed node");
})?;
+ Ok(())
+ }
+
+ async fn reset_pre_interop(&self) -> Result<(), ManagedNodeError> {
+ info!(target: "resetter", "Resetting the node to pre-interop state");
+ self.client.reset_pre_interop().await.inspect_err(|err| {
+ error!(target: "resetter", %err, "Failed to reset managed node to pre-interop state");
+ })?;
Ok(())
}
@@ -174,6 +187,7 @@ mod tests {
async fn pending_output_v0_at_timestamp(&self, timestamp: u64) -> Result;
async fn l2_block_ref_by_timestamp(&self, timestamp: u64) -> Result;
async fn block_ref_by_number(&self, block_number: u64) -> Result;
+ async fn reset_pre_interop(&self) -> Result<(), ClientError>;
async fn reset(&self, unsafe_id: BlockNumHash, cross_unsafe_id: BlockNumHash, local_safe_id: BlockNumHash, cross_safe_id: BlockNumHash, finalised_id: BlockNumHash) -> Result<(), ClientError>;
async fn provide_l1(&self, block_info: BlockInfo) -> Result<(), ClientError>;
async fn update_finalized(&self, finalized_block_id: BlockNumHash) -> Result<(), ClientError>;
@@ -214,7 +228,7 @@ mod tests {
#[tokio::test]
async fn test_reset_db_error() {
let mut db = MockDb::new();
- db.expect_get_super_head().returning(|| Err(StorageError::DatabaseNotInitialised));
+ db.expect_get_super_head().returning(|| Err(StorageError::LockPoisoned));
let client = MockClient::new();
diff --git a/crates/supervisor/core/src/syncnode/task.rs b/crates/supervisor/core/src/syncnode/task.rs
index 1c225c8c6e..2439a53283 100644
--- a/crates/supervisor/core/src/syncnode/task.rs
+++ b/crates/supervisor/core/src/syncnode/task.rs
@@ -295,6 +295,7 @@ mod tests {
async fn pending_output_v0_at_timestamp(&self, timestamp: u64) -> Result;
async fn l2_block_ref_by_timestamp(&self, timestamp: u64) -> Result;
async fn block_ref_by_number(&self, block_number: u64) -> Result;
+ async fn reset_pre_interop(&self) -> Result<(), ClientError>;
async fn reset(&self, unsafe_id: BlockNumHash, cross_unsafe_id: BlockNumHash, local_safe_id: BlockNumHash, cross_safe_id: BlockNumHash, finalised_id: BlockNumHash) -> Result<(), ClientError>;
async fn provide_l1(&self, block_info: BlockInfo) -> Result<(), ClientError>;
async fn update_finalized(&self, finalized_block_id: BlockNumHash) -> Result<(), ClientError>;
diff --git a/crates/supervisor/rpc/src/jsonrpsee.rs b/crates/supervisor/rpc/src/jsonrpsee.rs
index 693ee6223d..a2a3722ab8 100644
--- a/crates/supervisor/rpc/src/jsonrpsee.rs
+++ b/crates/supervisor/rpc/src/jsonrpsee.rs
@@ -172,6 +172,10 @@ pub trait ManagedModeApi {
#[method(name = "anchorPoint")]
async fn anchor_point(&self) -> RpcResult;
+ /// Reset the managed node to the pre-interop state
+ #[method(name = "resetPreInterop")]
+ async fn reset_pre_interop(&self) -> RpcResult<()>;
+
/// Reset the managed node to the specified block heads
#[method(name = "reset")]
async fn reset(
@@ -189,8 +193,8 @@ pub trait ManagedModeApi {
async fn fetch_receipts(&self, block_hash: BlockHash) -> RpcResult;
/// Get block infor for a given block number
- #[method(name = "blockRefByNumber")]
- async fn block_ref_by_number(&self, number: u64) -> RpcResult;
+ #[method(name = "l2BlockRefByNumber")]
+ async fn l2_block_ref_by_number(&self, number: u64) -> RpcResult;
/// Get the chain id
#[method(name = "chainID")]
diff --git a/crates/supervisor/storage/src/providers/head_ref_provider.rs b/crates/supervisor/storage/src/providers/head_ref_provider.rs
index 8e80dc7d3b..565c53544a 100644
--- a/crates/supervisor/storage/src/providers/head_ref_provider.rs
+++ b/crates/supervisor/storage/src/providers/head_ref_provider.rs
@@ -33,10 +33,7 @@ where
"Failed to seek head reference"
);
})?;
- let block_ref = result.ok_or_else(|| {
- warn!(target: "supervisor_storage", %safety_level, "No head reference found");
- StorageError::EntryNotFound("no head reference found".to_string())
- })?;
+ let block_ref = result.ok_or_else(|| StorageError::FutureData)?;
Ok(block_ref.into())
}
}
diff --git a/tests/devnets/simple-supervisor.yaml b/tests/devnets/simple-supervisor.yaml
index 5d265146ba..090fe86e95 100644
--- a/tests/devnets/simple-supervisor.yaml
+++ b/tests/devnets/simple-supervisor.yaml
@@ -9,7 +9,7 @@ optimism_package:
type: op-geth
cl:
type: op-node
- image: us-docker.pkg.dev/oplabs-tools-artifacts/images/op-node:v1.13.3
+ image: us-docker.pkg.dev/oplabs-tools-artifacts/images/op-node:develop
log_level: debug
network_params:
network: "kurtosis"
@@ -32,7 +32,7 @@ optimism_package:
type: op-geth
cl:
type: op-node
- image: us-docker.pkg.dev/oplabs-tools-artifacts/images/op-node:v1.13.3
+ image: us-docker.pkg.dev/oplabs-tools-artifacts/images/op-node:develop
log_level: debug
network_params:
network: "kurtosis"