use std::{ path::Path, sync::{Arc, Mutex}, time::Instant, }; use anyhow::{Context as _, Result, anyhow}; use borsh::BorshDeserialize; use chain_state::{ AcceptOutcome, Anchor, AnchorConsistencyCheck, ChainConsistency, ChainMismatch, ChainState, Tip, }; use common::{ HashType, block::{BedrockStatus, Block, BlockMeta, HashableBlockData}, transaction::{LeeTransaction, clock_invocation}, }; use config::{GenesisAction, SequencerConfig}; use futures::StreamExt as _; use itertools::Itertools as _; use lee::{AccountId, PublicTransaction, public_transaction::Message}; use lee_core::GENESIS_BLOCK_ID; use log::{error, info, warn}; use logos_blockchain_key_management_system_service::keys::{ED25519_SECRET_KEY_SIZE, Ed25519Key}; use logos_blockchain_zone_sdk::{ Slot, ZoneMessage, sequencer::{DepositInfo, WithdrawArg}, }; use mempool::{MemPool, MemPoolHandle}; #[cfg(feature = "mock")] pub use mock::SequencerCoreWithMockClients; use num_bigint::BigUint; pub use storage::error::DbError; use storage::sequencer::{ RocksDBIO, sequencer_cells::{PendingDepositEventRecord, WithdrawalReconciliationKey, ZoneAnchorRecord}, }; use crate::{ block_publisher::{BlockPublisherTrait, MsgId, ZoneSdkPublisher}, block_store::SequencerStore, }; pub mod block_publisher; pub mod block_store; pub mod config; pub mod cross_zone_watcher; #[cfg(feature = "mock")] pub mod mock; /// The origin of a transaction. #[derive(Clone, Copy)] pub enum TransactionOrigin { /// Basic transactions submitted by users via RPC. User, /// Transactions generated by the sequencer itself. Sequencer, } #[derive(Clone, Debug, BorshDeserialize)] struct DepositMetadata { recipient_id: lee::AccountId, } pub struct SequencerCore { /// Two-tier chain state: production builds on its head; the publisher's /// `on_follow` sink feeds adopted/orphaned/finalized peer blocks into it. chain: Arc>, store: SequencerStore, mempool: MemPool<(TransactionOrigin, LeeTransaction)>, sequencer_config: SequencerConfig, block_publisher: BP, } impl SequencerCore { /// Starts the sequencer using the provided configuration. /// If an existing database is found, the sequencer state is loaded from it and /// assumed to represent the correct latest state consistent with Bedrock-finalized data. /// If no database is found, the sequencer performs a fresh start from genesis, /// initializing its state with the accounts defined in the configuration file. fn open_or_create_store(config: &SequencerConfig) -> (SequencerStore, lee::V03State) { let signing_key = lee::PrivateKey::try_new(config.signing_key).unwrap(); let db_path = config.home.join("rocksdb"); if db_path.exists() { let store = SequencerStore::open_db(&db_path, signing_key).unwrap_or_else(|err| { panic!( "Failed to open database at {} with error: {err}", db_path.display() ) }); let state = store .get_lee_state() .expect("Failed to read state from store"); (store, state) } else { warn!( "Database not found at {}, starting from genesis", db_path.display() ); let (genesis_state, genesis_txs) = build_genesis_state(config); let hashable_data = HashableBlockData { block_id: GENESIS_BLOCK_ID, transactions: genesis_txs, prev_block_hash: HashType([0; 32]), timestamp: 0, }; let genesis_block = hashable_data.into_pending_block(&signing_key); let store = SequencerStore::create_db_with_genesis( &db_path, &genesis_block, &genesis_state, signing_key, ) .expect("Failed to create database with genesis block"); (store, genesis_state) } } /// Rebuilds the two-tier [`ChainState`]: the final tier from the persisted /// final snapshot (pre-genesis state when absent), the head tier by replaying /// every stored block above it, so a post-restart orphan can still revert. fn restore_chain_state( config: &SequencerConfig, store: &SequencerStore, stored_head_state: &lee::V03State, ) -> ChainState { let final_snapshot = store .dbio() .get_final_snapshot() .expect("Failed to read final snapshot from store"); let (final_state, final_tip) = match final_snapshot { Some((state, meta)) => (state, Some(Tip::from(meta))), // Nothing finalized yet: replay the whole stored chain. None => (build_initial_state(config), None), }; let boundary = final_tip.as_ref().map_or(0, |tip| tip.block_id); let mut head_blocks = store .get_all_blocks() .filter_ok(|block| block.header.block_id > boundary) .collect::, _>>() .expect("Failed to read blocks from store while restoring chain state"); head_blocks.sort_unstable_by_key(|block| block.header.block_id); let mut chain = ChainState::from_final(final_state, final_tip); for block in head_blocks { let block_id = block.header.block_id; chain.restore_head_block(block).unwrap_or_else(|err| { panic!("Stored block {block_id} does not replay while restoring chain state: {err}") }); } // The replayed head must reproduce the persisted state, else store // and config disagree (e.g. edited genesis actions). assert!( chain.head_state() == stored_head_state, "Persisted state does not match the replayed chain; reset the store or restore the original config" ); chain } pub async fn start_from_config( config: SequencerConfig, ) -> (Self, MemPoolHandle<(TransactionOrigin, LeeTransaction)>) { let bedrock_signing_key = load_or_create_signing_key(&config.home.join("bedrock_signing_key")) .expect("Failed to load or create bedrock signing key"); info!( "Bedrock signing public key: {}", hex::encode(bedrock_signing_key.public_key().to_bytes()) ); let (store, state) = Self::open_or_create_store(&config); let chain = Arc::new(Mutex::new(Self::restore_chain_state( &config, &store, &state, ))); let initial_checkpoint = store .get_zone_checkpoint() .expect("Failed to load zone-sdk checkpoint"); let is_fresh_start = initial_checkpoint.is_none(); let (mempool, mempool_handle) = MemPool::new(config.mempool_max_size); replay_unfulfilled_deposit_events(&store, mempool_handle.clone()); let block_publisher = BP::new( &config.bedrock_config, bedrock_signing_key, config.retry_pending_blocks_timeout, initial_checkpoint, Self::on_checkpoint(store.dbio()), Self::on_finalized_block(store.dbio()), Self::on_deposit_event(store.dbio(), mempool_handle.clone()), Self::on_withdraw_event(store.dbio()), Self::on_follow(store.dbio(), Arc::clone(&chain), mempool_handle.clone()), ) .await .expect("Failed to initialize Block Publisher"); // Cross-zone messaging: start a watcher per configured peer. The inbox // config account is seeded into genesis state in `build_genesis_state`. if let Some(cross_zone) = &config.cross_zone { cross_zone_watcher::spawn_watchers( &config.bedrock_config, cross_zone, config.block_create_timeout, &mempool_handle, ); } // Before producing, verify our local state still belongs to the chain // the channel serves and replay any channel blocks we are missing // (e.g. from other sequencers). let channel_absent = Self::verify_and_reconstruct(&block_publisher, &store, &chain, is_fresh_start) .await .expect("Failed to verify/reconstruct sequencer state from Bedrock"); // Publish our blocks only when we are bootstrapping a channel that does // not exist yet (no channel tip). If the channel already exists (another // sequencer created it), we adopted its blocks during reconstruction // instead; republishing then would fork the channel with our own copies. if is_fresh_start && channel_absent { let mut pending_blocks = store .get_all_blocks() .filter_ok(|block| matches!(block.bedrock_status, BedrockStatus::Pending)) .collect::, _>>() .expect("Failed to read blocks from store while republishing on fresh start"); pending_blocks.sort_unstable_by_key(|block| block.header.block_id); assert!( pending_blocks .first() .is_none_or(|block| block.header.block_id == GENESIS_BLOCK_ID), "First pending block on fresh start should be the genesis block" ); for block in &pending_blocks { block_publisher .publish_block(block, vec![]) .await .unwrap_or_else(|err| { panic!( "Failed to publish block {} on fresh start: {err:#}", block.header.block_id ) }); } } let sequencer_core = Self { chain, store, mempool, sequencer_config: config, block_publisher, }; (sequencer_core, mempool_handle) } /// Verifies the local store still belongs to the chain the connected channel /// serves and replays any finalized channel blocks missing locally into /// `state`/`store`, recording each block's L1 inscription slot as the new /// anchor. Fails (never parks) on any divergence. /// /// Returns whether the channel does not exist yet (has no tip), i.e. whether /// this sequencer is the one that must bootstrap-publish its own blocks. async fn verify_and_reconstruct( publisher: &BP, store: &SequencerStore, chain: &Mutex, is_fresh_start: bool, ) -> Result { let anchor_record = store .get_zone_anchor() .context("Failed to read zone anchor")?; let after_slot = anchor_record .and_then(|record| record.slot.checked_sub(1)) .map(Slot::from); let channel_tip_slot = publisher .channel_tip_slot() .await .context("Failed to read channel tip slot")?; // If this sequencer has already committed blocks to the channel, that // channel must still exist. A missing channel then means a wiped/rewound // Bedrock or a node pointing at a different chain, so refuse to resume // onto a foreign channel. // // "Committed" requires *both* a non-genesis tip and a checkpoint that was // persisted before this startup: the tip alone is set the moment we produce // (before the channel confirms it), while a checkpoint alone is written by // zone-sdk's cold-start backfill even on a brand-new empty channel before we // publish genesis. We must read the checkpoint presence from before `BP::new` // ran (`!is_fresh_start`), because its cold-start backfill re-persists a // checkpoint by the time we reach here — reading the store now would always // see one. Together they mean we produced blocks and zone-sdk processed // channel activity in a prior run. let local_tip = store .latest_block_meta() .context("Failed to read latest block meta")? .map(|meta| meta.id); let had_checkpoint_before_start = !is_fresh_start; if let Some(local_tip) = local_tip && had_checkpoint_before_start && channel_tip_slot.is_none() { return Err(anyhow!( "Sequencer holds committed blocks (tip {local_tip}) but the Bedrock channel \ no longer exists on the connected chain — the channel was wiped or the node \ points at a different chain. Refusing to resume onto a foreign channel." )); } let divergence_error = |mismatch: &ChainMismatch| { anyhow!( "Sequencer store diverges from the Bedrock channel ({mismatch}). \ Delete the sequencer storage directory or point at the correct channel." ) }; // With a recorded anchor, probe the channel for positive evidence of a // different chain: the frontier upfront (a missing/behind channel serves // no messages to scan), then the anchor block as messages stream in. let mut consistency_check = anchor_record.map(|record| { let anchor = Anchor::new( Slot::from(record.slot), Some((record.block_id, record.hash)), ); let mut check = AnchorConsistencyCheck::new(anchor); check.check_frontier(channel_tip_slot); check }); if let Some(ChainConsistency::Inconsistent(mismatch)) = consistency_check .as_ref() .and_then(AnchorConsistencyCheck::verdict) { return Err(divergence_error(mismatch)); } // Verify each message against the anchor and replay the // blocks (applying the ones we miss, checking the ones we hold). let messages = publisher .read_channel_after(after_slot) .await .context("Failed to read channel history for reconstruction")?; let mut messages = std::pin::pin!(messages); while let Some((message, slot)) = messages.next().await { if let Some(check) = &mut consistency_check && let Some(ChainConsistency::Inconsistent(mismatch)) = check.observe(&message, slot) { return Err(divergence_error(mismatch)); } let ZoneMessage::Block(zone_block) = message else { continue; }; let block: Block = borsh::from_slice(&zone_block.data).map_err(|err| { anyhow!( "Failed to deserialize channel block at slot {}: {err}", slot.into_inner() ) })?; // Locked per message (not across the stream `await`): concurrent // follow events interleave safely — both paths apply idempotently // and persist under this same lock. let mut chain = chain.lock().expect("chain state mutex poisoned"); Self::apply_reconstructed_block(store, &mut chain, &block, slot)?; } // The channel exists once it has a tip; only when it has none is this // sequencer the one bootstrapping it. This is deliberately not the // reconstruction scan's view above, which reads only finalized history // (up to LIB) and so reports "empty" while finality lags even though the // channel already holds unfinalized blocks from another sequencer. Ok(channel_tip_slot.is_none()) } /// Applies a single channel block during reconstruction: idempotent for /// blocks we already hold (verifying their hash), a validated continuation /// for new ones. Advances the persisted anchor to the block's slot. fn apply_reconstructed_block( store: &SequencerStore, chain: &mut ChainState, block: &Block, slot: Slot, ) -> Result<()> { let tip = store .latest_block_meta() .context("Failed to read latest block meta")?; let block_id = block.header.block_id; let block_hash = block.header.hash; let record = ZoneAnchorRecord { slot: slot.into_inner(), block_id, hash: block_hash, }; // A block at/below the tip must match what we already stored, otherwise // the channel is a different chain. if let Some(tip) = &tip && block_id <= tip.id { match store .get_block_at_id(block_id) .context("Failed to read stored block")? { Some(stored) if stored.header.hash == block_hash => { store .set_zone_anchor(&record) .context("Failed to persist zone anchor")?; return Ok(()); } Some(stored) => { return Err(anyhow!( "Channel block {block_id} hash {block_hash} does not match stored hash {}", stored.header.hash )); } None => { return Err(anyhow!( "Channel block {block_id} is at/below local tip {} but is missing locally", tip.id )); } } } // New continuation: channel history is finalized, so it goes through // the final tier — validation happens inside `apply_finalized`. match chain.apply_finalized(MsgId::from(block.header.hash.0), block, slot) { AcceptOutcome::Applied | AcceptOutcome::AlreadyApplied => {} AcceptOutcome::Parked(err) | AcceptOutcome::RetryableFailure(err) => { return Err(anyhow!( "Channel block {block_id} does not extend local tip {:?}: {err}", tip.map(|tip| tip.id) )); } } // Persist like the follow path: the tip meta stays pinned to the head // tip even when the reconstructed block lands below it. let head_tip = chain.head_tip().map(|head| BlockMeta::from(&head)); let final_meta = chain.final_tip().map(|meta| BlockMeta::from(&meta)); store .dbio() .store_followed_blocks( &[(block, true)], head_tip.as_ref(), chain.head_state(), final_meta.as_ref().map(|meta| (chain.final_state(), meta)), ) .context("Failed to persist reconstructed block")?; // Mark the deposits' pending records submitted so the production-time // guard drops the mints cold-start backfill re-queued for them. Withdraw // intents are deliberately not counted: backfill already re-delivered and // dropped their finalized L1 events, so an increment here would never be // consumed and would leave a phantom count. let deposit_event_ids: Vec<_> = block .body .transactions .iter() .filter_map(extract_bridge_deposit_id) .collect(); store .mark_deposit_events_submitted(&deposit_event_ids, block_id) .context("Failed to mark reconstructed deposits submitted")?; store .set_zone_anchor(&record) .context("Failed to persist zone anchor")?; Ok(()) } fn on_checkpoint(dbio: Arc) -> block_publisher::CheckpointSink { Box::new(move |cp| { let bytes = match serde_json::to_vec(&cp) { Ok(b) => b, Err(err) => { error!("Failed to serialize zone-sdk checkpoint: {err:#}"); return; } }; if let Err(err) = dbio.put_zone_sdk_checkpoint_bytes(&bytes) { error!("Failed to persist zone-sdk checkpoint: {err:#}"); } }) } fn on_finalized_block(dbio: Arc) -> block_publisher::FinalizedBlockSink { Box::new(move |block_id| { // NOTE: Theoretically Zone SDK may report finalization happening multiple times for the // same block. In practice this is very unlikely to happen. For that to // happen Sequencer should crash between receiving Finalized and Checkpoint events while // these events happen very fast (because Checkpoints are generated by Zone SDK // locally). if let Err(err) = dbio.clean_pending_blocks_up_to(block_id) { error!("Failed to mark pending blocks finalized up to {block_id}: {err:#}"); } match dbio.remove_fulfilled_pending_deposit_events_up_to_block(block_id) { Ok(0) => {} Ok(removed) => { info!( "Removed {removed} fulfilled pending deposit events up to finalized block {block_id}" ); } Err(err) => { error!( "Failed to remove fulfilled pending deposit events up to block {block_id}: {err:#}" ); } } }) } fn on_deposit_event( dbio: Arc, mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>, ) -> block_publisher::OnDepositEventSink { Box::new(move |deposit| { // NOTE: Theoretically Zone SDK may report multiple identical deposits. In practice this // is very unlikely to happen. For that to happen Sequencer should crash // between receiving Deposit and Checkpoint events while these events happen // very fast (because Checkpoints are generated by Zone SDK locally). let dbio = Arc::clone(&dbio); let mempool_handle = mempool_handle.clone(); Box::pin(async move { let id_hex = hex::encode(deposit.op_id); info!("Observed Bedrock Deposit event with id: {id_hex}"); let event_record = pending_deposit_event_record(&deposit); match dbio.add_pending_deposit_event(event_record.clone()) { Ok(true) => {} Ok(false) => { info!( "Deposit event {id_hex} already persisted as unfulfilled, skipping duplicate enqueue", ); return; } Err(err) => { error!( "Failed to persist unfulfilled deposit event {id_hex} before enqueue: {err:#}. Deposit will be lost.", ); return; } } let tx = match build_bridge_deposit_tx_from_event(&event_record) { Ok(tx) => tx, Err(err) => { error!( "Failed to build transaction from Bedrock deposit event {id_hex}: {err:#}. Deposit will be lost.", ); return; } }; if let Err(err) = mempool_handle .push((TransactionOrigin::Sequencer, tx)) .await { error!( "Failed to queue sequencer transaction built from finalized Bedrock event: {err:#}. Deposit will be lost." ); } }) }) } fn on_withdraw_event(dbio: Arc) -> block_publisher::OnWithdrawEventSink { Box::new(move |withdraw| { let dbio = Arc::clone(&dbio); Box::pin(async move { let hash_encoded = hex::encode(withdraw.tx_hash.as_ref()); let withdraw_key = match withdraw_event_reconciliation_key(&withdraw.op.outputs) { Ok(key) => key, Err(err) => { error!( "Failed to build reconciliation key for Bedrock Withdraw event with tx_hash {hash_encoded}: {err:#}" ); return; } }; match dbio.consume_unseen_withdraw_count(withdraw_key) { Ok(true) => { info!("Validated Bedrock Withdraw event with tx_hash: {hash_encoded}"); } Ok(false) => warn!( "Unexpected Bedrock Withdraw event with tx_hash {hash_encoded}: no matching unseen withdraw found" ), Err(err) => error!( "Failed to reconcile Bedrock Withdraw event with tx_hash {hash_encoded}: {err:#}" ), } }) }) } /// Publisher sink adapter over [`apply_follow_update`]. fn on_follow( dbio: Arc, chain: Arc>, mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>, ) -> block_publisher::OnFollowSink { Box::new(move |update: block_publisher::FollowUpdate| { apply_follow_update(&dbio, &chain, &mempool_handle, update); }) } /// Produces a new block from mempool transactions and publishes it via zone-sdk. pub async fn produce_new_block(&mut self) -> Result { let BlockWithMeta { block, deposit_event_ids, withdrawals, } = self .build_block_from_mempool() .context("Failed to build block from mempool transactions")?; let withdrawal_reconciliation_keys = withdrawals .iter() .map(|withdraw| withdraw_event_reconciliation_key(&withdraw.outputs)) .collect::>() .context("Failed to build reconciliation keys for block withdrawals")?; let this_msg = self .block_publisher .publish_block(&block, withdrawals) .await .context("Failed to publish block to Bedrock")?; self.record_produced_block( this_msg, &block, &deposit_event_ids, withdrawal_reconciliation_keys, )?; Ok(block.header.block_id) } /// Applies our own freshly-published block to the head with the [`MsgId`] the /// publish assigned it, so the head advances and the later adopted /// redelivery dedups, then persists it. /// /// Persistence is gated on the block actually becoming the head: if a peer /// block won this height while we were publishing (`AlreadyApplied`, or /// `Parked` when the head reorged to a different parent), the canonical /// block is persisted by the follow path instead, and our invalidated /// inscription comes back via `orphaned`. fn record_produced_block( &mut self, this_msg: MsgId, block: &Block, deposit_event_ids: &[HashType], withdrawal_reconciliation_keys: Vec, ) -> Result<()> { let mut chain = self.chain.lock().expect("chain state mutex poisoned"); match chain.apply_produced(this_msg, block) { AcceptOutcome::Applied => { // Persisted under the lock so disk writes land in apply order // with the follow path. self.store.update( block, deposit_event_ids, withdrawal_reconciliation_keys, chain.head_state(), )?; } AcceptOutcome::AlreadyApplied => { warn!( "Produced block {} lost a competing-write race, skipping persistence", block.header.block_id ); } AcceptOutcome::Parked(err) | AcceptOutcome::RetryableFailure(err) => { warn!( "Produced block {} no longer chains on the head, skipping persistence: {err}", block.header.block_id ); } } Ok(()) } /// Validates and applies a single mempool transaction to the current state. /// Returns `Ok(true)` if the transaction was valid and applied, `Ok(false)` if /// it was skipped due to validation failure. #[expect( clippy::too_many_arguments, reason = "splitting the produce-path accumulators into a struct buys nothing" )] fn apply_mempool_transaction( store: &SequencerStore, state: &mut lee::V03State, origin: TransactionOrigin, tx: &LeeTransaction, block_height: u64, timestamp: u64, deposit_event_ids: &mut Vec, withdrawals: &mut Vec, ) -> Result { let tx_hash = tx.hash(); match origin { TransactionOrigin::User => { let validated_diff = match tx.validate_on_state(state, block_height, timestamp) { Ok(diff) => diff, Err(err) => { error!( "Transaction with hash {tx_hash} failed execution check with error: {err:#?}, skipping it", ); return Ok(false); } }; if let Some(withdraw_data) = extract_bridge_withdraw_data(tx) { withdrawals.push(withdraw_data); } state.apply_state_diff(validated_diff); } TransactionOrigin::Sequencer => { let LeeTransaction::Public(public_tx) = tx else { panic!("Sequencer may only generate Public transactions, found {tx:#?}"); }; if let Some(deposit_op_id) = extract_bridge_deposit_id(tx) { if store .is_deposit_event_submitted(deposit_op_id) .context("Failed to check whether deposit was already submitted")? { info!("Skipping already-submitted bridge deposit {deposit_op_id}"); return Ok(false); } deposit_event_ids.push(deposit_op_id); } state .transition_from_public_transaction(public_tx, block_height, timestamp) .context("Failed to execute sequencer-generated transaction")?; } } info!("Validated transaction with hash {tx_hash}, including it in block"); Ok(true) } fn build_block_from_mempool(&mut self) -> Result { let now = Instant::now(); // Build on the head: its tip is the parent, its state the validation base. let (prev_block_hash, new_block_height, mut working_state) = { let chain = self.chain.lock().expect("chain state mutex poisoned"); let tip = chain.head_tip(); let height = tip.as_ref().map_or(GENESIS_BLOCK_ID, |head| { head.block_id .checked_add(1) .expect("block id should not overflow") }); let prev = tip.map_or(HashType([0; 32]), |head| head.hash); (prev, height, chain.head_state().clone()) }; let mut valid_transactions = Vec::new(); let mut deposit_event_ids = Vec::new(); let mut withdrawals = Vec::new(); let max_block_size = usize::try_from(self.sequencer_config.max_block_size.as_u64()) .expect("`max_block_size` should fit into usize"); let new_block_timestamp = u64::try_from(chrono::Utc::now().timestamp_millis()) .expect("Timestamp must be positive"); let clock_tx = clock_invocation(new_block_timestamp); let clock_lee_tx = LeeTransaction::Public(clock_tx.clone()); while let Some((origin, tx)) = self.mempool.pop() { let tx_hash = tx.hash(); let temp_valid_transactions = [ valid_transactions.as_slice(), std::slice::from_ref(&tx), std::slice::from_ref(&clock_lee_tx), ] .concat(); let temp_hashable_data = HashableBlockData { block_id: new_block_height, transactions: temp_valid_transactions, prev_block_hash, timestamp: new_block_timestamp, }; let block_size = borsh::to_vec(&temp_hashable_data) .context("Failed to serialize block for size check")? .len(); if block_size > max_block_size { warn!( "Transaction with hash {tx_hash} deferred to next block: \ block size {block_size} bytes would exceed limit of {max_block_size} bytes", ); self.mempool.push_front((origin, tx)); break; } if Self::apply_mempool_transaction( &self.store, &mut working_state, origin, &tx, new_block_height, new_block_timestamp, &mut deposit_event_ids, &mut withdrawals, )? { valid_transactions.push(tx); } if valid_transactions.len() >= self.sequencer_config.max_num_tx_in_block { break; } } working_state .transition_from_public_transaction(&clock_tx, new_block_height, new_block_timestamp) .context("Clock transaction failed. Aborting block production.")?; valid_transactions.push(clock_lee_tx); let hashable_data = HashableBlockData { block_id: new_block_height, transactions: valid_transactions, prev_block_hash, timestamp: new_block_timestamp, }; let block = hashable_data .clone() .into_pending_block(self.store.signing_key()); log::info!( "Created block with {} transactions in {} seconds", hashable_data.transactions.len(), now.elapsed().as_secs() ); Ok(BlockWithMeta { block, deposit_event_ids, withdrawals, }) } /// Reads the current head state under the lock without cloning it, so callers /// reuse `V03State`'s own API (accounts, nonces, proofs) with no whole-state copy. pub fn with_state(&self, f: impl FnOnce(&lee::V03State) -> R) -> R { f(self .chain .lock() .expect("chain state mutex poisoned") .head_state()) } pub const fn block_store(&self) -> &SequencerStore { &self.store } #[must_use] pub fn chain_height(&self) -> u64 { self.chain .lock() .expect("chain state mutex poisoned") .head_tip() .map_or(0, |tip| tip.block_id) } pub const fn sequencer_config(&self) -> &SequencerConfig { &self.sequencer_config } /// Marks all pending blocks with `block_id <= last_finalized_block_id` as /// finalized. Idempotent. Production callers don't invoke this directly — /// it's wired up in `start_from_config` to the publisher's /// `on_finalized_block` sink, which fires on `Event::TxsFinalized` / /// `Event::FinalizedInscriptions`. Kept on the type for tests. // TODO: Delete blocks instead of marking them as finalized. Current // approach is used because we still have `GetBlockDataRequest`. pub fn clean_finalized_blocks_from_db(&self, last_finalized_block_id: u64) -> Result<()> { info!("Clearing pending blocks up to id: {last_finalized_block_id}"); self.store .dbio() .clean_pending_blocks_up_to(last_finalized_block_id)?; Ok(()) } /// Returns the list of stored pending blocks. pub fn get_pending_blocks(&self) -> Result> { Ok(self .store .get_all_blocks() .collect::>>()? .into_iter() .filter(|block| matches!(block.bedrock_status, BedrockStatus::Pending)) .collect()) } pub const fn block_publisher(&self) -> &BP { &self.block_publisher } /// Whether this sequencer is currently authorized to write to the channel. #[must_use] pub fn is_our_turn(&self) -> bool { self.block_publisher.is_our_turn() } /// Shared handle to the two-tier follow state, for tests to drive the /// follow path directly. #[cfg(all(test, feature = "mock"))] fn chain(&self) -> Arc> { Arc::clone(&self.chain) } } struct BlockWithMeta { block: Block, deposit_event_ids: Vec, withdrawals: Vec, } /// Feed one channel delta into the follow state and mirror it to the store: /// revert orphaned, then apply and persist adopted and finalized blocks. /// Production builds on this same head. Wired to the publisher via /// [`SequencerCore::on_follow`]; a free function so tests can drive it directly. /// /// TODO: unlike the indexer's ingest loop, this path does not retry /// `is_retryable` (transient) apply failures — a failed block just parks and /// relies on a valid successor or a restart. `ChainState` never emits /// `AcceptOutcome::RetryableFailure` yet; adding retry parity here is a /// follow-up. fn apply_follow_update( dbio: &RocksDBIO, chain: &Mutex, mempool_handle: &MemPoolHandle<(TransactionOrigin, LeeTransaction)>, update: block_publisher::FollowUpdate, ) { let block_publisher::FollowUpdate { adopted, orphaned, finalized, } = update; // The lock is held across the persist below so disk writes land in apply // order — the produce path persists under this same lock. let resubmit_txs = { let mut chain = chain.lock().expect("chain state mutex poisoned"); // User txs of orphaned blocks, returned to the mempool below. let resubmit_txs: Vec = orphaned .iter() .flat_map(|(_, block)| resubmittable_txs(block)) .collect(); // Outcomes align with `adopted`. let outcomes = chain.apply_channel_update(&orphaned, &adopted); let mut to_persist: Vec<(&Block, bool)> = adopted .iter() .zip(&outcomes) .filter(|(_, outcome)| matches!(outcome, AcceptOutcome::Applied)) .map(|((_, block), _)| (block, false)) .collect(); let mut final_advanced = false; for (this_msg, block) in &finalized { // FIXME: thread the finalized inscription's L1 slot once the // sdk surfaces it; only used for the invalid-finalized stall. if matches!( chain.apply_finalized(*this_msg, block, Slot::from(0)), AcceptOutcome::Applied ) { to_persist.push((block, true)); final_advanced = true; } } // Snapshot the advanced final tier so a restart re-anchors on it. let final_meta = final_advanced.then(|| { let tip = chain.final_tip().expect("advanced final tier has a tip"); BlockMeta::from(&tip) }); let head_tip = chain.head_tip().map(|tip| BlockMeta::from(&tip)); // One atomic write for the whole update: blocks, tip meta and the // state after the last block land together, so a crash can never // leave the stored state ahead of the stored blocks. A persist // failure is fatal: the in-memory chain has already advanced, and // continuing would leave a permanent gap in the store. // // The `panic!` ends the drive task, whose cancellation halts the node. // // TODO: the zone-sdk checkpoint is persisted by `on_checkpoint` // before this write; a crash in between resumes past these blocks // without them ever landing in the store. Full `BlocksProcessed` // atomicity (checkpoint + blocks + state in one batch, per the sdk's // event contract) is a follow-up. dbio.store_followed_blocks( &to_persist, head_tip.as_ref(), chain.head_state(), final_meta.as_ref().map(|meta| (chain.final_state(), meta)), ) .unwrap_or_else(|err| panic!("Failed to persist followed blocks: {err:#}")); resubmit_txs }; // Rebuild orphaned work: return its user txs to the mempool so the // next on-turn production re-includes them on the new head. // // We use [`try_push`] here because this is called from the // publisher's drive task, and only the block production drains the mempool. // A blocking push on a full mempool would deadlock here. for tx in resubmit_txs { let tx_hash = tx.hash(); if let Err(err) = mempool_handle.try_push((TransactionOrigin::User, tx)) { warn!("Dropping orphaned transaction {tx_hash} on resubmit: {err}"); } } } /// Checks the database for any pending deposit events that have not yet been marked as submitted in /// a block, and re-queues them in the mempool in a separate async task for inclusion in the next /// block. fn replay_unfulfilled_deposit_events( store: &SequencerStore, mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>, ) { let replay_records: Vec = store .get_unfulfilled_deposit_events() .expect("Failed to load unfulfilled deposit events") .into_iter() .filter(|record| record.submitted_in_block_id.is_none()) .collect(); if replay_records.is_empty() { return; } info!( "Found {} unfulfilled deposit events in DB, re-queueing", replay_records.len() ); tokio::spawn(async move { for record in replay_records { let tx = match build_bridge_deposit_tx_from_event(&record) { Ok(tx) => tx, Err(err) => { warn!( "Skipping replay of pending deposit event {} due to tx build failure: {err:#}", hex::encode(record.deposit_op_id) ); continue; } }; if let Err(err) = mempool_handle .push((TransactionOrigin::Sequencer, tx)) .await { error!( "Failed to re-queue unfulfilled deposit event {} from DB: {err:#}", hex::encode(record.deposit_op_id) ); break; } } }); } /// The pre-genesis state: `testnet_initial_state` plus the bridge-lock holdings, /// the only accounts seeded outside any transaction. Cross-zone config is seeded /// by genesis `InitConfig` transactions and reconstructed by replaying them. fn build_initial_state(config: &SequencerConfig) -> lee::V03State { #[cfg(not(feature = "testnet"))] let base = testnet_initial_state::initial_state(); #[cfg(feature = "testnet")] let base = testnet_initial_state::initial_state_testnet(); // Bridge-lock holder balances belong to the source side and are not produced by // any transaction, so seed them directly. Cross-zone config is seeded by genesis // InitConfig transactions in `build_genesis_state`, not here. let holdings = bridge_lock_holdings(&config.genesis) .map(|(holder, amount)| cross_zone::build_holding_account(holder, amount)); base.with_public_accounts(holdings) } /// Builds the initial genesis state from [`build_initial_state`] plus configured /// genesis transactions. Returns the final state and the list of /// [`LeeTransaction`]s that should be committed to the genesis block so external /// observers can replay them. fn build_genesis_state(config: &SequencerConfig) -> (lee::V03State, Vec) { let mut state = build_initial_state(config); // Fingerprint the directly-seeded state, before genesis txs, so it matches the indexer's. info!( "Genesis fingerprint: {}", hex::encode(state.genesis_fingerprint()) ); // Config txs seed the config accounts by transaction, so every node // reconstructs them by replaying the genesis block. The wrapped-token minter is // initialized on every zone (wrapped_token is a builtin), since its InitConfig // is user-callable and a config PDA left default would be claimable by anyone as // the first initializer (a minter hijack). The inbox allowlist is initialized // only on receiving zones; the inbox is sequencer-only, so its default config // PDA is not user-claimable, merely unused until the zone receives. let wrapped_token_config_tx = std::iter::once(cross_zone::build_wrapped_token_init_config_tx()); let inbox_config_tx = config.cross_zone.as_ref().map(|cross_zone| { let self_zone = *config.bedrock_config.channel_id.as_ref(); cross_zone::build_inbox_init_config_tx(self_zone, cross_zone) }); let supply_txs = config.genesis.iter().filter_map(|action| match action { GenesisAction::SupplyAccount { account_id, balance, } => Some(build_supply_account_genesis_transaction( account_id, *balance, )), GenesisAction::SupplyBridgeAccount { balance } => { Some(build_supply_bridge_account_genesis_transaction(*balance)) } // Seeded directly in `build_initial_state` (holdings via `build_holding_account`), not a // genesis tx. GenesisAction::SupplyBridgeLockHolding { .. } => None, }); let genesis_txs = wrapped_token_config_tx .chain(inbox_config_tx) .chain(supply_txs) .chain(std::iter::once(clock_invocation(0))) .inspect(|tx| { state .transition_from_public_transaction(tx, GENESIS_BLOCK_ID, 0) .expect("Failed to execute genesis transaction"); }) .map(LeeTransaction::Public) .collect(); (state, genesis_txs) } /// Bridge-lock holder balances configured for this zone's genesis. fn bridge_lock_holdings( genesis: &[GenesisAction], ) -> impl Iterator + '_ { genesis.iter().filter_map(|action| match action { GenesisAction::SupplyBridgeLockHolding { holder, amount } => Some((*holder, *amount)), GenesisAction::SupplyAccount { .. } | GenesisAction::SupplyBridgeAccount { .. } => None, }) } /// Whether a program may only be invoked by sequencer-origin transactions. /// /// The cross-zone inbox is injected solely by the watcher; a user-submitted call /// must be rejected at ingress, since `TransactionOrigin` is not carried in the /// block. #[must_use] pub fn is_sequencer_only_program(program_id: lee::ProgramId) -> bool { cross_zone::is_sequencer_only_program(program_id) } fn build_supply_account_genesis_transaction( account_id: &AccountId, balance: u128, ) -> PublicTransaction { let faucet_program_id = programs::faucet().id(); let vault_program_id = programs::vault().id(); let recipient_vault_id = vault_core::compute_vault_account_id(vault_program_id, *account_id); let message = Message::try_new( faucet_program_id, vec![system_accounts::faucet_account_id(), recipient_vault_id], vec![], faucet_core::Instruction::GenesisTransferVault { vault_program_id, recipient_id: *account_id, amount: balance, }, ) .expect("Failed to serialize genesis transfer instruction"); let witness_set = lee::public_transaction::WitnessSet::from_raw_parts(vec![]); PublicTransaction::new(message, witness_set) } fn build_supply_bridge_account_genesis_transaction(balance: u128) -> PublicTransaction { let faucet_program_id = programs::faucet().id(); let bridge_account_id = system_accounts::bridge_account_id(); let message = Message::try_new( faucet_program_id, vec![system_accounts::faucet_account_id(), bridge_account_id], vec![], faucet_core::Instruction::GenesisTransferDirect { amount: balance }, ) .expect("Failed to serialize bridge genesis transfer instruction"); let witness_set = lee::public_transaction::WitnessSet::from_raw_parts(vec![]); PublicTransaction::new(message, witness_set) } fn pending_deposit_event_record(deposit: &DepositInfo) -> PendingDepositEventRecord { PendingDepositEventRecord { deposit_op_id: HashType(deposit.op_id), source_tx_hash: HashType(deposit.tx_hash.0), amount: deposit.amount, metadata: deposit.metadata.clone().into(), submitted_in_block_id: None, } } fn build_bridge_deposit_tx_from_event(event: &PendingDepositEventRecord) -> Result { let metadata = DepositMetadata::try_from_slice(&event.metadata) .context("Failed to decode finalized Bedrock deposit metadata")?; let bridge_program_id = programs::bridge().id(); let vault_program_id = programs::vault().id(); let recipient_vault_id = vault_core::compute_vault_account_id(vault_program_id, metadata.recipient_id); let message = Message::try_new( bridge_program_id, vec![system_accounts::bridge_account_id(), recipient_vault_id], vec![], bridge_core::Instruction::Deposit { l1_deposit_op_id: event.deposit_op_id.0, vault_program_id, recipient_id: metadata.recipient_id, amount: event.amount, }, ) .context("Failed to build bridge deposit message")?; let witness_set = lee::public_transaction::WitnessSet::from_raw_parts(vec![]); Ok(LeeTransaction::Public(PublicTransaction::new( message, witness_set, ))) } /// User transactions of an orphaned block to return to the mempool: everything /// except the trailing clock tx, sequencer-generated bridge deposits (replayed /// from their own bedrock events) and sequencer-only cross-zone txs (replayed /// by the watcher; the ingress guard rejects them as `User`). fn resubmittable_txs(block: &Block) -> Vec { let Some((_clock, rest)) = block.body.transactions.split_last() else { return Vec::new(); }; rest.iter() .filter(|tx| extract_bridge_deposit_id(tx).is_none() && !is_sequencer_only_tx(tx)) .cloned() .collect() } #[must_use] fn is_sequencer_only_tx(tx: &LeeTransaction) -> bool { matches!(tx, LeeTransaction::Public(tx) if is_sequencer_only_program(tx.message().program_id)) } #[must_use] fn extract_bridge_deposit_id(tx: &LeeTransaction) -> Option { let LeeTransaction::Public(tx) = tx else { return None; }; let message = tx.message(); if message.program_id != programs::bridge().id() { return None; } let instruction = risc0_zkvm::serde::from_slice::(&message.instruction_data) .ok()?; match instruction { bridge_core::Instruction::Deposit { l1_deposit_op_id, .. } => Some(HashType(l1_deposit_op_id)), bridge_core::Instruction::Withdraw { .. } => None, } } #[must_use] fn extract_bridge_withdraw_data(tx: &LeeTransaction) -> Option { let LeeTransaction::Public(tx) = tx else { return None; }; let message = tx.message(); if message.program_id != programs::bridge().id() { return None; } let instruction = risc0_zkvm::serde::from_slice::(&message.instruction_data) .ok()?; let bridge_core::Instruction::Withdraw { amount, bedrock_account_pk, } = instruction else { return None; }; let recipient_pk = logos_blockchain_key_management_system_service::keys::ZkPublicKey::from( BigUint::from_bytes_le(&bedrock_account_pk), ); Some(WithdrawArg { outputs: logos_blockchain_core::mantle::ledger::Outputs::new( logos_blockchain_core::mantle::Note::new(amount, recipient_pk), ), }) } fn withdraw_event_reconciliation_key( outputs: &logos_blockchain_core::mantle::ledger::Outputs, ) -> Result { let [note] = outputs.as_ref().as_slice() else { return Err(anyhow!( "Unsupported withdraw output count for reconciliation: {}", outputs.len() )); }; // `extract_bridge_withdraw_data` maps [u8;32] LE -> BigUint -> ZkPublicKey. // Reconcile by reversing that direction here. let mut bedrock_account_pk = BigUint::from(note.pk.into_inner()).to_bytes_le(); if bedrock_account_pk.len() > 32 { return Err(anyhow!( "Withdraw recipient public key is too large: {} bytes", bedrock_account_pk.len() )); } bedrock_account_pk.resize(32, 0); let bedrock_account_pk: [u8; 32] = bedrock_account_pk .try_into() .expect("Public key bytes were padded/truncated to 32 bytes"); Ok(WithdrawalReconciliationKey { amount: note.value, bedrock_account_pk, }) } /// Load signing key from file or generate a new one if it doesn't exist. pub fn load_or_create_signing_key(path: &Path) -> Result { if path.exists() { let key_bytes = std::fs::read(path)?; let key_array: [u8; ED25519_SECRET_KEY_SIZE] = key_bytes .try_into() .map_err(|_bytes| anyhow!("Found key with incorrect length"))?; Ok(Ed25519Key::from_bytes(&key_array)) } else { let mut key_bytes = [0_u8; ED25519_SECRET_KEY_SIZE]; rand::RngCore::fill_bytes(&mut rand::thread_rng(), &mut key_bytes); // Create parent directory if it doesn't exist if let Some(parent) = path.parent() { std::fs::create_dir_all(parent)?; } std::fs::write(path, key_bytes)?; Ok(Ed25519Key::from_bytes(&key_bytes)) } } #[cfg(test)] #[cfg(feature = "mock")] mod tests;