Files
logos-execution-zone/lez/sequencer/core/src/lib.rs
T

2387 lines
98 KiB
Rust

use std::{
collections::VecDeque,
path::Path,
sync::{Arc, Mutex},
time::{Duration, 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 cross_zone_inbox_core::CrossZoneMessage;
use futures::StreamExt as _;
use itertools::Itertools as _;
use lee::{AccountId, PublicTransaction, public_transaction::Message};
use lee_core::GENESIS_BLOCK_ID;
use log::{debug, error, info, warn};
use logos_blockchain_core::mantle::ops::channel::Ed25519PublicKey;
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;
// Re-exported because `cross_zone_dead_letters` returns it and the service
// crate does not depend on `storage`, so it could not otherwise name the type.
pub use storage::sequencer::sequencer_cells::DeadLetterDispatchRecord;
use storage::sequencer::{
DispatchFailure, RocksDBIO, StoreUpdate,
sequencer_cells::{
DispatchOrigin, PendingCrossZoneDispatchRecord, PendingDepositEventRecord,
WithdrawalReconciliationKey, ZoneAnchorRecord,
},
};
use tokio_retry::{Retry, strategy::FixedInterval};
use crate::{
block_publisher::{BlockPublisherTrait, MsgId, NoteId, ZoneSdkPublisher},
block_store::SequencerStore,
task_group::TaskGroup,
};
pub mod block_publisher;
pub mod block_store;
pub mod committee_discovery;
pub mod config;
pub mod cross_zone_watcher;
pub mod gossip;
#[cfg(feature = "mock")]
pub mod mock;
pub mod task_group;
/// Failed production attempts before a cross-zone dispatch is given up on.
///
/// One attempt per block, so this is tens of seconds of retrying. Enough for a
/// failure that is not the message's fault to clear, short enough that a message
/// which will never execute stops being retried.
const RETIRE_DISPATCH_AFTER_FAILURES: u32 = 3;
/// Cross-zone deliveries one block may carry.
///
/// Each one costs a guest execution whether it succeeds or fails, and what
/// queues them up is chosen by peer zones. Without a bound, a backlog decides
/// how long a block takes to build and leaves no room for user transactions,
/// since store-drained work is taken before the mempool. The rest wait one
/// block; nothing is dropped.
const MAX_DISPATCHES_PER_BLOCK: usize = 16;
/// Fixed, public key behind a genesis-only funding account: the faucet can
/// only be called top-level, not as `Stake`'s mover, so this account is a
/// pass-through that receives faucet funds and then moves them into the real
/// stake account. Not a secret: every node derives the same account, and it
/// holds nothing once genesis has run.
// TODO: replace the faucet pass-through with a real deposit from Bedrock,
// once that path exists, instead of a fixed genesis-only key.
const GENESIS_STAKE_FUNDING_KEY: [u8; 32] = [9; 32];
/// A number of Bedrock slots, as opposed to a [`Slot`] position.
type SlotCount = u64;
/// A founding sequencer's key, plus the ownership account attesting to its stake.
type FoundingStake = (
sequencer_stake_core::SequencerKey,
lee::PublicKey,
lee::Signature,
);
/// 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,
/// Transactions received via p2p gossip from a peer sequencer.
Gossip,
}
impl From<TransactionOrigin> for sequencer_core_metrics::TransactionOrigin {
fn from(origin: TransactionOrigin) -> Self {
match origin {
TransactionOrigin::User => Self::User,
TransactionOrigin::Sequencer => Self::Sequencer,
TransactionOrigin::Gossip => Self::Gossip,
}
}
}
#[derive(Clone, Debug, BorshDeserialize)]
struct DepositMetadata {
recipient_id: lee::AccountId,
}
pub struct SequencerCore<BP: BlockPublisherTrait = ZoneSdkPublisher> {
/// 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<Mutex<ChainState>>,
store: SequencerStore,
mempool: MemPool<(TransactionOrigin, LeeTransaction)>,
sequencer_config: SequencerConfig,
block_publisher: BP,
/// Cross-zone watchers, stopped when this sequencer is dropped. They hold a
/// store handle, so leaving them running would keep the `RocksDB` lock held
/// and make the home directory unopenable by a restarting sequencer.
watchers: TaskGroup,
/// Channel tip slot as of the last committee-config submission.
last_committee_submission_slot: Option<Slot>,
}
impl<BP: BlockPublisherTrait> SequencerCore<BP> {
const CHANNEL_PROBE_RETRIES: usize = 29;
const CHANNEL_PROBE_RETRY_DELAY: Duration = Duration::from_secs(2);
/// Channel slots between committee-config submissions; a margin over
/// observed Bedrock confirmation lag.
const COMMITTEE_SUBMISSION_COOLDOWN: SlotCount = 10;
/// 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,
bootstrap_sequencer_key: Option<sequencer_stake_core::SequencerKey>,
) -> (SequencerStore, lee::V03State) {
let signing_key = lee::PrivateKey::try_new(config.signing_key).unwrap();
let db_path = config.db_path();
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 {
let legacy = config.home.join("rocksdb");
if legacy.exists() {
warn!(
"Ignoring pre-channel-suffix database at {}; rename it to {} to resume it",
legacy.display(),
db_path.display()
);
}
warn!(
"Database not found at {}, starting from genesis",
db_path.display()
);
let (genesis_state, genesis_txs) = build_genesis_state(config, bootstrap_sequencer_key);
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");
// Incrementing count for genesis.
sequencer_core_metrics::increment_blocks_produced_total();
(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::<Result<Vec<_>, _>>()
.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)>) {
sequencer_core_metrics::init();
let bedrock_signing_key =
load_or_create_signing_key(&config.home.join("bedrock_signing_key"))
.expect("Failed to load or create bedrock signing key");
log::info!(
"Bedrock signing public key: {}",
hex::encode(bedrock_signing_key.public_key().to_bytes())
);
let own_sequencer_key =
sequencer_stake_core::SequencerKey::new(bedrock_signing_key.public_key().to_bytes())
.expect("our own Bedrock public key is a valid Ed25519 public key");
// Only seed our own key into genesis as the bootstrap sequencer if the
// channel doesn't exist yet. Otherwise it's someone else's channel and
// we join later, the normal self-join way.
let channel_probe_retry_strategy =
FixedInterval::new(Self::CHANNEL_PROBE_RETRY_DELAY).take(Self::CHANNEL_PROBE_RETRIES);
let channel_already_exists = Retry::start(channel_probe_retry_strategy, || async {
BP::channel_exists(&config.bedrock_config)
.await
.inspect_err(|err| warn!("Failed to probe Bedrock channel: {err:#}"))
})
.await
.expect("Failed to probe Bedrock channel");
if channel_already_exists {
info!("Channel already exists; joining as a non channel creator");
} else {
info!("Channel does not exist yet; starting it as channel creator");
}
let bootstrap_sequencer_key = (!channel_already_exists).then_some(own_sequencer_key);
let (store, state) = Self::open_or_create_store(&config, bootstrap_sequencer_key);
assert!(
committee_discovery::config_is_readable(&state),
"sequencer_stake config account is absent or undecodable; this chain's state is not \
one this sequencer can operate on"
);
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);
sequencer_core_metrics::record_mempool_max_size(config.mempool_max_size);
let block_publisher = BP::new(
&config.bedrock_config,
bedrock_signing_key,
config.retry_pending_blocks_timeout,
initial_checkpoint,
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`.
let watchers = config
.cross_zone
.as_ref()
.map_or_else(TaskGroup::default, |cross_zone| {
cross_zone_watcher::spawn_watchers(
&config.bedrock_config,
cross_zone,
config.block_create_timeout,
&store.dbio(),
)
});
// 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");
// Seed the high water mark from the tip we are starting on. Every stored
// block reached the store by being published or by being adopted from
// the channel, so the channel holds them all and none is ours to write
// again. Without this the mark is absent until the first publish of this
// run, leaving that window unguarded — which is exactly the window a
// store written before the mark existed starts in.
if let Some(tip) = store
.latest_block_meta()
.expect("Failed to read latest block meta")
{
store
.raise_published_high_water(tip.id)
.expect("Failed to seed published high water mark");
}
// 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::<Result<Vec<_>, _>>()
.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"
);
// The channel is born holding only its creator's key, so a configured
// founding set is applied by the same tx that writes genesis; the
// committee is never observable without it.
let founding_committee = founding_committee(&config, own_sequencer_key);
let mut last_checkpoint = None;
for block in &pending_blocks {
let publish = match &founding_committee {
Some(keys) if block.header.block_id == GENESIS_BLOCK_ID => {
block_publisher
.publish_genesis_creating_channel(block, keys.clone())
.await
}
_ => block_publisher.publish_block(block, vec![]).await,
};
let outcome = publish.unwrap_or_else(|err| {
panic!(
"Failed to publish block {} on fresh start: {err:#}",
block.header.block_id
)
});
last_checkpoint = Some(outcome.checkpoint);
store
.raise_published_high_water(block.header.block_id)
.expect("Failed to persist published high water mark");
}
// These blocks are already stored, so only the sdk's pending set
// moved. Checkpoints are cumulative — persisting just the last one
// is both sufficient and the only way to keep this loop linear.
if let Some(checkpoint) = last_checkpoint {
store
.set_zone_checkpoint(&checkpoint)
.expect("Failed to persist checkpoint after republishing on fresh start");
}
}
let sequencer_core = Self {
chain,
store,
mempool,
sequencer_config: config,
block_publisher,
watchers,
last_committee_submission_slot: None,
};
sequencer_core_metrics::record_chain_height(sequencer_core.chain_height());
record_dead_letter_gauge(&sequencer_core.store.dbio());
(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) when the channel proves a different chain:
/// the anchor consistency check, or a block that will not validate.
///
/// 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<ChainState>,
is_fresh_start: bool,
) -> Result<bool> {
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, ignored when it conflicts at a height the final
/// tier already settled, a validated continuation otherwise. 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 we already hold verbatim needs no replay, but the channel
// serving it is what makes it irreversible, so its deliveries are
// settled and their records are owed nothing. Without this a restart
// leaves a record for every delivery it already published, and nothing
// downstream would ever remove them.
if let Some(tip) = &tip
&& block_id <= tip.id
&& let Some(stored) = store
.get_block_at_id(block_id)
.context("Failed to read stored block")?
&& stored.header.hash == block_hash
{
settle_reconstructed_deliveries(store, &stored);
store
.set_zone_anchor(&record)
.context("Failed to persist zone anchor")?;
return Ok(());
}
// A conflict at a height the final tier already settled: the channel
// carries two inscriptions for one block id — competing sequencers
// around a turn change — and finality already picked one, so the other
// is dropped. `apply_adopted` ignores the same conflict. A genuinely
// foreign channel is caught upstream by the anchor consistency check,
// not here; the anchor stays on the block we hold.
if let Some(final_tip) = chain.final_tip()
&& block_id <= final_tip.block_id
{
log::warn!(
"Ignoring channel block {block_id} with hash {block_hash} conflicting with the \
finalized block at this height"
);
return Ok(());
}
// Above the final tier the head is reorg-able, so finalized history
// wins: `apply_finalized` finalizes the matching prefix and rebases the
// head onto what the channel settled. Validation happens inside it.
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)
));
}
}
// A reconstructed block is finalized, so any deposit it mints is
// permanently reflected in state (its receipt PDA); drop the pending
// record backfill may have re-delivered, so the drain stops re-minting.
let finalized_deposit_ids: Vec<_> = block
.body
.transactions
.iter()
.filter_map(extract_bridge_deposit_id)
.collect();
// The same for the deliveries it carries: the inbox has seen them, so
// the drain would skip them anyway, and the records are owed nothing.
let finalized_dispatch_keys = settled_dispatch_keys(&store.dbio(), block);
// The tip meta stays pinned to the head tip even when the reconstructed
// block lands below it, and the anchor only advances if the block
// itself landed.
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_update(&StoreUpdate {
blocks: &[(block, true)],
head_tip: head_tip.as_ref(),
final_snapshot: final_meta.as_ref().map(|meta| (chain.final_state(), meta)),
remove_deposit_records: &finalized_deposit_ids,
remove_dispatch_records: &finalized_dispatch_keys,
zone_anchor: Some(&record),
..StoreUpdate::new(chain.head_state())
})
.context("Failed to persist reconstructed block")?;
Ok(())
}
/// Publisher sink adapter over [`apply_follow_update`].
fn on_follow(
dbio: Arc<RocksDBIO>,
chain: Arc<Mutex<ChainState>>,
mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>,
) -> block_publisher::OnFollowSink {
Box::new(move |update: block_publisher::FollowUpdate| {
apply_follow_update(&dbio, &chain, &mempool_handle, update);
})
}
/// Runs everything this sequencer owes its turn: builds a block from
/// mempool transactions, publishes it via zone-sdk, and submits any
/// committee-config update the new state calls for.
pub async fn run_production_turn(&mut self) -> Result<u64> {
let live_accredited_keys = self.live_accredited_sequencer_keys().await;
let BlockWithMeta {
block,
withdrawals,
committee_update,
} = self
.build_block_from_mempool(live_accredited_keys.as_deref())
.context("Failed to build block from mempool transactions")?;
let block_publisher::PublishOutcome {
this_msg,
checkpoint,
released_notes,
} = self
.block_publisher
.publish_block(&block, withdrawals)
.await
.context("Failed to publish block to Bedrock")?;
// The inscription is on L1 from here on, whatever the head does with the
// block below, so this height must never be published again.
self.store
.raise_published_high_water(block.header.block_id)
.context("Failed to persist published high water mark")?;
// Independent Mantle tx, not bundled with the block above — join/exit
// config updates don't need to be.
self.submit_committee_update(committee_update).await;
let withdrawal_reconciliation_keys: Vec<_> = released_notes
.iter()
.map(withdrawal_reconciliation_key)
.collect();
self.record_produced_block(
this_msg,
&block,
&withdrawal_reconciliation_keys,
&checkpoint,
)?;
Ok(block.header.block_id)
}
/// Live committee snapshot for gating `FinalizeUnstake` inclusion and
/// committee updates. `None` if it could not be read.
async fn live_accredited_sequencer_keys(
&self,
) -> Option<Vec<sequencer_stake_core::SequencerKey>> {
match self.block_publisher.accredited_keys().await {
Ok(keys) => Some(
keys.iter()
.filter_map(|key| {
sequencer_stake_core::SequencerKey::new(key.to_bytes()).or_else(|| {
warn!(
"Ignoring accredited key {}: not a valid Ed25519 public key",
hex::encode(key.to_bytes())
);
None
})
})
.collect(),
),
Err(err) => {
warn!(
"Failed to read live committee snapshot; skipping FinalizeUnstake inclusion \
and committee updates this round: {err:#}"
);
None
}
}
}
/// Whether the channel has advanced far enough past `last_submission` to
/// submit again. A missing tip counts as no advance.
fn committee_cooldown_elapsed(last_submission: Option<Slot>, tip: Option<Slot>) -> bool {
let Some(last_submission) = last_submission else {
return true;
};
tip.is_some_and(|tip| {
tip.into_inner()
.saturating_sub(last_submission.into_inner())
>= Self::COMMITTEE_SUBMISSION_COOLDOWN
})
}
async fn submit_committee_update(
&mut self,
committee_update: Option<Vec<sequencer_stake_core::SequencerKey>>,
) {
let Some(new_keys) = committee_update else {
return;
};
let tip_slot = match self.block_publisher.channel_tip_slot().await {
Ok(tip_slot) => tip_slot,
Err(err) => {
warn!("Failed to read channel tip slot; skipping committee update: {err:#}");
return;
}
};
if !Self::committee_cooldown_elapsed(self.last_committee_submission_slot, tip_slot) {
return;
}
let new_keys = new_keys
.into_iter()
.map(|key| {
Ed25519PublicKey::from_bytes(&key.to_bytes())
.expect("sequencer key was decoded from a valid Ed25519 public key")
})
.collect();
self.last_committee_submission_slot = tip_slot;
if let Err(err) = self.block_publisher.submit_channel_config(new_keys).await {
warn!("Failed to submit committee channel-config update: {err:#}");
}
}
/// 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,
withdrawal_reconciliation_keys: &[WithdrawalReconciliationKey],
checkpoint: &block_publisher::SequencerCheckpoint,
) -> Result<()> {
let checkpoint_bytes = block_store::checkpoint_bytes(checkpoint)?;
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,
withdrawal_reconciliation_keys,
chain.head_state(),
Some(&checkpoint_bytes),
)?;
sequencer_core_metrics::increment_blocks_produced_total();
sequencer_core_metrics::record_chain_height(block.header.block_id);
}
// Neither branch persists anything, checkpoint included: the
// inscription it holds as pending belongs to a block that is not
// ours to keep.
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.
fn apply_mempool_transaction(
state: &mut lee::V03State,
origin: TransactionOrigin,
tx: &LeeTransaction,
block_height: u64,
timestamp: u64,
withdrawals: &mut Vec<WithdrawArg>,
) -> bool {
let tx_hash = tx.hash();
match origin {
// Gossiped transactions arrive from untrusted peers, same as
// user-submitted ones, so they get the same full state validation.
TransactionOrigin::User | TransactionOrigin::Gossip => {
let validated_diff = match tx.validate_on_state(state, block_height, timestamp) {
Ok(diff) => diff,
Err(err) => {
// A gossiped tx the leader already included is
// expected to fail here (e.g. on nonce) for every
// other node on its turn; that is steady-state noise,
// not an error. User-submitted failures still warrant
// `error!`.
if matches!(origin, TransactionOrigin::Gossip) {
debug!(
"Transaction with hash {tx_hash} failed execution check with error: {err:#?}, skipping it",
);
} else {
error!(
"Transaction with hash {tx_hash} failed execution check with error: {err:#?}, skipping it",
);
}
return 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:#?}");
};
// Bridge deposits are deduped by their receipt PDA in chain
// state (drained only when unminted, no-op on replay), so no
// node-local guard is needed here.
//
// Skip-and-log rather than propagate: a drained deposit is
// re-fed from the store every turn and only finality removes it,
// so a `?` here would let a single unexecutable mint (e.g. a
// bridge escrow under-funded relative to the L1 deposit, which
// every sequencer hits identically) abort production on all of
// them forever. Skipping keeps the record queued for retry
// without halting the node.
if let Err(err) =
state.transition_from_public_transaction(public_tx, block_height, timestamp)
{
error!(
"Sequencer-generated transaction {tx_hash} failed execution: {err:#?}, skipping it",
);
return false;
}
}
}
log::info!("Validated transaction with hash {tx_hash}, including it in block");
true
}
fn build_block_from_mempool(
&mut self,
live_accredited_keys: Option<&[sequencer_stake_core::SequencerKey]>,
) -> Result<BlockWithMeta> {
let now = Instant::now();
// Decoded outside the chain lock, and read before it is taken: the usual
// case is no delivery records at all, and decoding is the expensive part.
// One that does not decode is dropped rather than kept, since nothing
// will ever turn those bytes into a block transaction.
let mut settled = Vec::new();
let recorded_dispatches: Vec<_> = self
.store
.dbio()
.get_pending_cross_zone_dispatches()
.context("Failed to load pending cross-zone dispatches")?
.into_iter()
.filter_map(
|record| match borsh::from_slice::<LeeTransaction>(&record.transaction) {
Ok(tx) => {
let message = extract_cross_zone_dispatch(&tx);
Some((record.message_key, message, tx))
}
Err(err) => {
warn!(
"Dropping pending cross-zone dispatch {} that does not decode: {err:#}",
hex::encode(record.message_key)
);
settled.push(record.message_key);
None
}
},
)
.collect();
// Build on the head: its tip is the parent, its state the validation
// base.
//
// The delivery records are classified in here rather than after, so the
// final state can be read by reference. Cloning it cost a full state
// copy on every block of every zone, cross-zone or not.
let (
prev_block_hash,
new_block_height,
mut working_state,
pending_dispatches,
finalize_unstake_txs,
) = {
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);
// Three outcomes per record. Already in the final state means the
// delivery is irreversible, so the record is dropped; that is the
// only thing that removes a record the watcher re-added after its
// delivery had already settled, which it does whenever it re-reads a
// slot it has consumed. Already in the head state but not the final
// one means the delivery is on this chain but could still orphan, so
// the record is skipped and kept. Otherwise it goes in this block.
let mut pending: VecDeque<LeeTransaction> = VecDeque::new();
for (key, message, tx) in recorded_dispatches {
match message {
Some(message) if dispatch_already_delivered(chain.final_state(), &message) => {
settled.push(key);
}
Some(message) if dispatch_already_delivered(chain.head_state(), &message) => {}
_ if pending.len() >= MAX_DISPATCHES_PER_BLOCK => {}
_ => pending.push_back(tx),
}
}
(
prev,
height,
chain.head_state().clone(),
pending,
build_finalize_unstake_txs(chain.head_state()),
)
};
if !settled.is_empty() {
if let Err(err) = self
.store
.dbio()
.drop_settled_cross_zone_dispatches(&settled)
{
// Only bookkeeping: the deliveries themselves are irreversible,
// and the next turn tries again.
warn!(
"Failed to drop {} settled delivery record(s): {err:#}",
settled.len()
);
}
// A settled delivery may be one this node had given up on, which
// takes its dead letter with it.
record_dead_letter_gauge(&self.store.dbio());
}
let mut valid_transactions = Vec::new();
let mut withdrawals = Vec::new();
// Bridge deposit mints are drained from the store, not the mempool: the
// follow path records the event durably but cannot enqueue the mint
// itself (it runs on the publisher's drive task, where an await stalls
// the very task production needs). Draining here also subsumes the old
// startup replay.
//
// Skip any deposit whose receipt PDA already exists in the state we
// build on — it was minted by us or by a peer whose block we adopted.
// An orphan reverts the receipt with the block, so the next turn
// re-mints without any bookkeeping of our own.
let pending_deposits: VecDeque<LeeTransaction> = self
.store
.get_pending_deposit_events()
.context("Failed to load pending deposit events")?
.into_iter()
.filter(|record| !deposit_already_minted(&working_state, record.deposit_op_id))
.filter_map(|record| {
build_bridge_deposit_tx_from_event(&record)
.inspect_err(|err| {
warn!(
"Skipping pending deposit event {} due to tx build failure: {err:#}",
hex::encode(record.deposit_op_id)
);
})
.ok()
})
.collect();
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());
sequencer_core_metrics::record_mempool_size(self.mempool.len());
// Everything drained from the store first, then user work. `from_store`
// is not the same as a `Sequencer` origin: it says the transaction has a
// record behind it and so needs no requeue, where the origin only says
// it was not submitted by a user.
let mut pending_from_store = pending_deposits;
pending_from_store.extend(pending_dispatches);
pending_from_store.extend(finalize_unstake_txs);
while let Some((origin, tx, from_store)) = pending_from_store
.pop_front()
.map(|tx| (TransactionOrigin::Sequencer, tx, true))
.or_else(|| self.mempool.pop().map(|(origin, tx)| (origin, tx, false)))
{
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 {
// Would a block carrying nothing but this still be too big? Then
// it does not fit in any block and deferring it defers it for
// ever. A store-drained transaction is at the head of the queue
// every turn, so breaking here would stop production reaching
// anything behind it, including the whole mempool, permanently.
// Count it against the delivery instead so it is given up on.
//
// Measured on its own rather than from `block_size`, which also
// counts whatever this block already holds: a transaction that
// merely does not fit *today* is the ordinary deferral below.
if from_store
&& !self.fits_in_an_empty_block(
&tx,
&clock_lee_tx,
new_block_height,
prev_block_hash,
new_block_timestamp,
)?
{
error!(
"Sequencer-drained transaction {tx_hash} cannot fit in any block under the \
{max_block_size} byte limit; giving up on it rather than stalling production",
);
self.count_dispatch_failure(&tx);
continue;
}
warn!(
"Transaction with hash {tx_hash} deferred to next block: \
block size {block_size} bytes would exceed limit of {max_block_size} bytes",
);
// Anything drained from the store needs no requeue: its record
// stays there and is drained again on the next turn.
if !from_store {
self.mempool.push_front((origin, tx));
}
break;
}
// Block-validity rule: a not-yet-valid FinalizeUnstake is dropped
// outright, not applied — whether it arrived via the mempool
// (anyone may submit one, per spec) or from this sequencer's own
// discovery above. It re-appears on its own once conditions are
// met (mempool: whoever wants it finalized resubmits;
// discovery-sourced: reconstructed fresh next block), so it
// doesn't need requeuing here.
if !finalize_unstake_is_includable(&working_state, &tx, live_accredited_keys) {
continue;
}
let before_tx_apply = Instant::now();
let applied = Self::apply_mempool_transaction(
&mut working_state,
origin,
&tx,
new_block_height,
new_block_timestamp,
&mut withdrawals,
);
if applied {
sequencer_core_metrics::record_mempool_transaction_application_time(
origin.into(),
tx.kind().into(),
sequencer_core_metrics::ApplyStatus::Applied,
before_tx_apply.elapsed(),
);
valid_transactions.push(tx);
} else {
sequencer_core_metrics::increment_mempool_failed_transactions_total();
sequencer_core_metrics::record_mempool_transaction_application_time(
origin.into(),
tx.kind().into(),
sequencer_core_metrics::ApplyStatus::Failed,
before_tx_apply.elapsed(),
);
// A failed transaction is simply left out of the block, except a
// dispatch: that one is re-fed from the store every turn, so one
// that can never execute would fail on every block for ever.
self.count_dispatch_failure(&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);
sequencer_core_metrics::record_transactions_per_block(valid_transactions.len());
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()
);
sequencer_core_metrics::record_block_creation_time(now.elapsed());
let committee_update = live_accredited_keys
.and_then(|keys| committee_discovery::committee_update(&working_state, keys));
Ok(BlockWithMeta {
block,
withdrawals,
committee_update,
})
}
/// 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<R>(&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 no longer calls this: finalization
/// flips now ride the follow path's atomic write via
/// [`StoreUpdate::finalized_up_to`]. 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<()> {
log::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<Vec<Block>> {
Ok(self
.store
.get_all_blocks()
.collect::<block_store::DbResult<Vec<Block>>>()?
.into_iter()
.filter(|block| matches!(block.bedrock_status, BedrockStatus::Pending))
.collect())
}
pub const fn block_publisher(&self) -> &BP {
&self.block_publisher
}
/// Whether a block carrying nothing but `tx` and the clock would be within
/// the size limit.
///
/// Distinguishes "does not fit in this block" from "does not fit in any
/// block". The first is an ordinary deferral; the second, for a transaction
/// the store re-feeds every turn, is a permanent stall unless it is given up
/// on.
fn fits_in_an_empty_block(
&self,
tx: &LeeTransaction,
clock_tx: &LeeTransaction,
block_id: u64,
prev_block_hash: HashType,
timestamp: u64,
) -> Result<bool> {
let alone = HashableBlockData {
block_id,
transactions: vec![tx.clone(), clock_tx.clone()],
prev_block_hash,
timestamp,
};
let size = borsh::to_vec(&alone)
.context("Failed to serialize block for size check")?
.len();
let max = usize::try_from(self.sequencer_config.max_block_size.as_u64())
.expect("`max_block_size` should fit into usize");
Ok(size <= max)
}
/// Counts one failed production attempt against `tx` if it is a cross-zone
/// delivery, giving up on it once too many accumulate.
///
/// A delivery's payload and target accounts are chosen on the peer zone and
/// validated by nobody in between, so one can fail for good; but a failure
/// can equally be a property of the moment, so give up only after several.
/// Giving up moves the record to the dead letter: a peer cannot grow the
/// pending list with deliveries that never execute, and it stays findable.
fn count_dispatch_failure(&self, tx: &LeeTransaction) {
let Some(message) = extract_cross_zone_dispatch(tx) else {
return;
};
let key = cross_zone_inbox_core::message_key(
&message.src_zone,
message.src_block_id,
message.src_tx_index,
);
let origin = DispatchOrigin {
src_zone: message.src_zone,
src_block_id: message.src_block_id,
src_tx_index: message.src_tx_index,
};
match self
.store
.dbio()
.record_dispatch_failure(key, RETIRE_DISPATCH_AFTER_FAILURES, origin)
{
Ok(DispatchFailure::Retired(record)) => {
sequencer_core_metrics::increment_cross_zone_dispatches_retired_total();
record_dead_letter_gauge(&self.store.dbio());
error!(
"Giving up on cross-zone delivery {} from peer zone {} block {} transaction {} ({} bytes) after {} failed attempts. This node will not retry it; unless another sequencer carries it, the message is not delivered. Kept in the dead letter.",
hex::encode(key),
hex::encode(origin.src_zone),
origin.src_block_id,
origin.src_tx_index,
record.transaction_bytes,
record.failed_attempts
);
}
Ok(DispatchFailure::Retried { failed_attempts }) => warn!(
"Cross-zone delivery {} failed to execute ({failed_attempts} of {RETIRE_DISPATCH_AFTER_FAILURES} attempts), will retry next block",
hex::encode(key)
),
// Not a give-up: the ordinary case is a delivery that already
// settled, so its record is gone and there is nothing left to lose.
Ok(DispatchFailure::Absent) => debug!(
"Cross-zone delivery {} failed to execute but has no pending record; nothing to count",
hex::encode(key)
),
Err(err) => error!(
"Failed to count the failed attempt for cross-zone delivery {}: {err:#}",
hex::encode(key)
),
}
}
/// The deliveries this node has given up on, and how many times it has.
///
/// Retained is read first so the pair can only skew towards a total that
/// leads its list, an ordinary evicted or settled state. The other order
/// would report entries against a total of zero.
pub fn cross_zone_dead_letters(&self) -> Result<(u64, Vec<DeadLetterDispatchRecord>), DbError> {
let dbio = self.store.dbio();
let retained = dbio.get_dead_letter_cross_zone_dispatches()?;
let total = dbio.get_dead_letter_cross_zone_dispatch_count()?;
Ok((total, retained))
}
/// Every background task that holds this sequencer's store handle.
///
/// Taken before the core is shared, so a shutdown path can wait for them
/// without owning the core. Until all of them have stopped the `RocksDB`
/// lock is still held and the home directory cannot be reopened, which is
/// what a restart does.
#[must_use]
pub fn background_tasks(&self) -> Vec<TaskGroup> {
vec![
self.watchers.clone(),
self.block_publisher.background_tasks(),
]
}
/// 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()
}
/// The height the next produced block would claim.
#[must_use]
pub fn next_block_height(&self) -> u64 {
self.chain
.lock()
.expect("chain state mutex poisoned")
.head_tip()
.map_or(GENESIS_BLOCK_ID, |tip| {
tip.block_id
.checked_add(1)
.expect("block id should not overflow")
})
}
/// `Some(high_water)` when the head has rewound below what we already
/// inscribed, so the next block would be a *second*, different block at a
/// height the channel already carries. Callers must skip their turn.
///
/// The head alone cannot detect this: an orphan report for our own
/// still-unfinalized blocks rewinds it (`ChainState::apply_channel_update`)
/// and prunes those blocks from the store, so the tip reads as if they were
/// never produced. The mark is kept outside that pruning for exactly this.
///
/// This is not a stall — the head recovers by itself once the inscriptions
/// we are protecting finalize and the final tier rebases onto them.
#[must_use]
pub fn rewound_below_published(&self) -> Option<u64> {
let high_water = self.store.published_high_water().ok().flatten()?;
(self.next_block_height() <= high_water).then_some(high_water)
}
/// 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<Mutex<ChainState>> {
Arc::clone(&self.chain)
}
}
struct BlockWithMeta {
block: Block,
withdrawals: Vec<WithdrawArg>,
committee_update: Option<Vec<sequencer_stake_core::SequencerKey>>,
}
/// Whether `deposit_op_id`'s mint is already reflected in `state` — its receipt
/// PDA exists. The receipt is the exactly-once ledger the bridge program keeps.
fn deposit_already_minted(state: &lee::V03State, deposit_op_id: HashType) -> bool {
let receipt_id =
bridge_core::deposit_receipt_account_id(programs::bridge().id(), deposit_op_id.0);
state
.get_account_by_id_ref(receipt_id)
.is_some_and(|receipt| *receipt != lee::Account::default())
}
/// Whether a cross-zone delivery is already on the chain we are building on.
///
/// The inbox records each peer block's delivered indices in that block's seen
/// shard and no-ops a replay, so the shard is the same kind of answer the
/// deposit receipt gives: state, not bookkeeping. An orphan reverts the entry
/// with the block, so the next turn re-delivers with nothing to unwind.
///
/// Both halves matter. A shard bound to a different peer block is not this
/// delivery's replay record, it is what will make it abort, and calling that
/// delivered would drop the record instead of dead-lettering it.
fn dispatch_already_delivered(state: &lee::V03State, message: &CrossZoneMessage) -> bool {
let shard_id = cross_zone_inbox_core::inbox_seen_shard_account_id(
programs::cross_zone_inbox().id(),
&message.src_zone,
message.src_block_id,
);
state.get_account_by_id_ref(shard_id).is_some_and(|shard| {
cross_zone_inbox_core::SeenShard::from_bytes(shard.data.as_ref()).is_ok_and(|seen| {
seen.binds(&message.src_block_hash) && seen.contains(message.src_tx_index)
})
})
}
/// Publishes how many given-up-on deliveries are retained.
///
/// Read from the store because the list falls as well as rises (eviction, and
/// reconciliation when a delivery settles elsewhere). Costs a read and a decode,
/// so call it only where one of those can have happened.
fn record_dead_letter_gauge(dbio: &RocksDBIO) {
match dbio.get_dead_letter_cross_zone_dispatches() {
Ok(records) => {
sequencer_core_metrics::record_cross_zone_dead_letter_dispatches(records.len());
}
Err(err) => {
warn!("Failed to read the cross-zone dead letter for its gauge: {err:#}");
}
}
}
/// 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.
///
/// Everything the event produced lands in one write — see [`StoreUpdate`].
///
/// 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<ChainState>,
mempool_handle: &MemPoolHandle<(TransactionOrigin, LeeTransaction)>,
update: block_publisher::FollowUpdate,
) {
let block_publisher::FollowUpdate {
checkpoint,
adopted,
orphaned,
finalized,
deposits,
withdrawals,
} = update;
let checkpoint_bytes = block_store::checkpoint_bytes(&checkpoint)
.unwrap_or_else(|err| panic!("Failed to serialize zone-sdk checkpoint: {err:#}"));
// NOTE: Theoretically Zone SDK may re-deliver an already seen deposit or
// finalization. Both are idempotent here: a deposit already on record is
// not re-appended, and a finalization only ever moves the tier forward.
let deposit_records: Vec<PendingDepositEventRecord> =
deposits.iter().map(pending_deposit_event_record).collect();
// One reconciliation unit per released note, matching how the intents were
// recorded at publish time.
let consumed_withdrawals: Vec<WithdrawalReconciliationKey> = withdrawals
.iter()
.flat_map(|withdraw| withdraw.op.inputs.iter())
.map(withdrawal_reconciliation_key)
.collect();
// 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, outcome, head_height) = {
let mut chain = chain.lock().expect("chain state mutex poisoned");
// An orphan report rewinds the head to the earliest orphaned block and
// prunes the store above it, which is how a run of our own inscriptions
// can silently stop being ours. Loud on the way in: it is the only
// trace, and the rewind it causes is the expensive one.
// Debug, not warn: the sdk orphans our blocks routinely once LIB pruning
// drops them from the lineage, and most of those no longer sit in the
// head. The rewind below is the part that costs something.
let head_before = chain.head_tip().map(|tip| tip.block_id);
if !orphaned.is_empty() {
let ids: Vec<u64> = orphaned
.iter()
.map(|(_, block)| block.header.block_id)
.collect();
debug!(
"Channel orphaned {} block(s) {:?}..={:?}, head tip is {head_before:?}",
ids.len(),
ids.iter().min(),
ids.iter().max(),
);
}
// Outcomes align with `adopted`.
let outcomes = chain.apply_channel_update(&orphaned, &adopted);
// An adoption that does not apply freezes the head where it is, and
// every later one then fails the same way. Nothing else reports it.
for ((_, block), outcome) in adopted.iter().zip(&outcomes) {
if let AcceptOutcome::Parked(err) | AcceptOutcome::RetryableFailure(err) = outcome {
warn!(
"Adopted block {} did not apply, head stays at {:?}: {err}",
block.header.block_id,
chain.head_tip().map(|tip| tip.block_id),
);
}
}
if let (Some(before), Some(after)) = (head_before, chain.head_tip().map(|tip| tip.block_id))
&& after < before
{
warn!("Head rewound from {before} to {after}");
}
let mut to_persist: Vec<(&Block, bool)> = adopted
.iter()
.zip(&outcomes)
.filter(|(_, outcome)| matches!(outcome, AcceptOutcome::Applied))
.map(|((_, block), _)| (block, false))
.collect();
// Only blocks the final tier holds drive the bookkeeping below: a parked
// one never became irreversible, so marking blocks finalized through it
// or dropping its deposit records would lose them for good.
let mut irreversible: Vec<&Block> = Vec::new();
let mut final_advanced = false;
for (this_msg, block) in &finalized {
// FIXME: thread the finalized inscription's L1 slot instead of
// `Slot::from(0)`; only used for the invalid-finalized stall.
// logos-blockchain PR #3147 surfaces it as `FinalizedTx.l1_slot` —
// wire it through `FollowUpdate::finalized` once the zone-sdk pin is
// bumped past that (a separate PR).
match chain.apply_finalized(*this_msg, block, Slot::from(0)) {
AcceptOutcome::Applied => {
to_persist.push((block, true));
irreversible.push(block);
final_advanced = true;
}
// A re-delivery of a block the final tier already holds: no new
// payload and the tier does not move, but it is irreversible all
// the same, so it still settles its deposits.
AcceptOutcome::AlreadyApplied => irreversible.push(block),
AcceptOutcome::Parked(_) | AcceptOutcome::RetryableFailure(_) => {}
}
}
// User txs of orphaned blocks, returned to the mempool below.
//
// Computed after the finalized tier has advanced, and only for blocks
// above it: the zone-sdk reports a block as orphaned once LIB pruning
// drops its inscription from the channel lineage, so every block of
// ours is orphaned a poll or two after it finalizes. Those transactions
// are irreversibly included, and returning them to the mempool puts
// them back in every block we produce from then on.
let final_height = chain.final_tip().map(|tip| tip.block_id);
let resubmit_txs: Vec<LeeTransaction> = orphaned
.iter()
.filter(|(_, block)| final_height.is_none_or(|id| block.header.block_id > id))
.flat_map(|(_, block)| resubmittable_txs(block))
.collect();
// 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));
// Every block at or below the highest finalized one is irreversible, so
// stored blocks there can be marked finalized.
let last_finalized = irreversible.iter().map(|block| block.header.block_id).max();
// A deposit observed in a finalized block is permanently minted (its
// receipt is now in the irreversible tier), so its pending record can be
// dropped. Keyed by op id, not block id: a record only goes once its own
// deposit finalizes, never because some other block finalized at its
// height.
let finalized_deposit_ids: Vec<HashType> = irreversible
.iter()
.flat_map(|block| block.body.transactions.iter())
.filter_map(extract_bridge_deposit_id)
.collect();
// The same for cross-zone deliveries, keyed by message key: a record
// goes once its own delivery is irreversible, never because another
// block finalized at its height.
let finalized_dispatch_keys: Vec<[u8; 32]> = irreversible
.iter()
.flat_map(|block| settled_dispatch_keys(dbio, block))
.collect();
// 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.
let outcome = dbio
.store_update(&StoreUpdate {
checkpoint: Some(&checkpoint_bytes),
blocks: &to_persist,
head_tip: head_tip.as_ref(),
final_snapshot: final_meta.as_ref().map(|meta| (chain.final_state(), meta)),
finalized_up_to: last_finalized,
new_deposit_events: &deposit_records,
remove_deposit_records: &finalized_deposit_ids,
remove_dispatch_records: &finalized_dispatch_keys,
consumed_withdrawals: &consumed_withdrawals,
..StoreUpdate::new(chain.head_state())
})
.unwrap_or_else(|err| panic!("Failed to persist follow update: {err:#}"));
(resubmit_txs, outcome, head_tip.map_or(0, |tip| tip.id))
};
sequencer_core_metrics::record_chain_height(head_height);
// The runtime reconcile path: finalizing another sequencer's block drops the
// dead letter of a delivery this node gave up on.
record_dead_letter_gauge(dbio);
if outcome.accepted_deposits > 0 {
log::info!(
"Recorded {} Bedrock Deposit event(s); their mints are drained from the store on our next turn",
outcome.accepted_deposits
);
}
for withdrawal in &outcome.unmatched_withdrawals {
warn!(
"Unexpected Bedrock Withdraw event releasing channel note {}: no matching unseen withdraw found",
hex::encode(withdrawal.released_note_id)
);
}
// 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 block production drains the mempool. A blocking
// push would stall the drive task, and a sequencer that is not on turn
// never produces — so nothing would ever drain it again.
//
// TODO: a full mempool still drops the transaction; a durable resubmit
// queue is a follow-up.
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}");
}
}
}
/// The pre-genesis state: `testnet_initial_state` plus the bridge-lock holdings,
/// the only accounts seeded outside any transaction. Everything else, including
/// the bootstrap sequencer's own stake, is applied as a genesis transaction in
/// [`build_genesis_state`] so followers replay it instead of guessing it.
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,
bootstrap_sequencer_key: Option<sequencer_stake_core::SequencerKey>,
) -> (lee::V03State, Vec<LeeTransaction>) {
let mut state = build_initial_state(config);
// Fingerprint the directly-seeded state, before genesis txs, so it matches the indexer's.
log::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 and
// both emitters' pins are initialized on every zone: all three are builtins with
// a user-callable InitConfig, so a config PDA left default is claimable by the
// first initializer, hijacking the minter or repointing an emitter's outbox. 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(
config.cross_zone.as_ref(),
));
let ping_sender_config_tx = std::iter::once(cross_zone::build_ping_sender_init_config_tx());
let ping_receiver_config_tx = std::iter::once(cross_zone::build_ping_receiver_init_config_tx(
config.cross_zone.as_ref(),
));
let bridge_lock_config_tx = std::iter::once(cross_zone::build_bridge_lock_init_config_tx());
let inbox_config_tx = config.cross_zone.as_ref().map(|_| {
let self_zone = *config.bedrock_config.channel_id.as_ref();
cross_zone::build_inbox_init_config_tx(self_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. Stakes are built separately below.
GenesisAction::SupplyBridgeLockHolding { .. } | GenesisAction::StakeSequencer { .. } => {
None
}
});
// The creator falls back to staking itself, signing with the key it owns.
let mut staked = founding_stakes(config);
if staked.is_empty() {
staked.extend(bootstrap_sequencer_key.map(|key| {
let key_path = config.home.join("sequencer_stake_signing_key");
let owner = load_or_create_stake_signing_key(&key_path)
.expect("Failed to load or create the stake signing key");
let signature = sign_genesis_stake(0, key, &owner);
(key, lee::PublicKey::new_from_private_key(&owner), signature)
}));
}
let bootstrap_stake_txs = build_stake_genesis_transactions(&staked);
let genesis_txs = wrapped_token_config_tx
.chain(ping_sender_config_tx)
.chain(ping_receiver_config_tx)
.chain(bridge_lock_config_tx)
.chain(inbox_config_tx)
.chain(supply_txs)
.chain(bootstrap_stake_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)
}
fn founding_stakes(config: &SequencerConfig) -> Vec<FoundingStake> {
config
.genesis
.iter()
.filter_map(|action| match action {
GenesisAction::StakeSequencer {
sequencer_key,
ownership_public_key,
stake_signature,
} => Some((
*sequencer_key,
ownership_public_key.clone(),
stake_signature.clone(),
)),
GenesisAction::SupplyAccount { .. }
| GenesisAction::SupplyBridgeAccount { .. }
| GenesisAction::SupplyBridgeLockHolding { .. } => None,
})
.collect()
}
/// The accredited keys a newly created channel should carry, `own_key` first
/// because creation gives the turn to index 0. `None` leaves creation to the
/// plain inscription path.
fn founding_committee(
config: &SequencerConfig,
own_key: sequencer_stake_core::SequencerKey,
) -> Option<Vec<block_publisher::Ed25519PublicKey>> {
let mut keys: Vec<_> = founding_stakes(config)
.into_iter()
.map(|(key, ..)| key)
.collect();
if keys.is_empty() {
return None;
}
keys.sort_unstable();
keys.retain(|key| *key != own_key);
Some(
std::iter::once(own_key)
.chain(keys)
.map(|key| {
block_publisher::Ed25519PublicKey::from_bytes(&key.to_bytes())
.expect("sequencer key was decoded from a valid Ed25519 public key")
})
.collect(),
)
}
fn genesis_stake_funding_account() -> AccountId {
let key = lee::PrivateKey::try_new(GENESIS_STAKE_FUNDING_KEY)
.expect("GENESIS_STAKE_FUNDING_KEY is a valid private key");
AccountId::from(&lee::PublicKey::new_from_private_key(&key))
}
/// The exact `Stake` message the founding sequencer at `index` must sign. Shared
/// offchain by the genesis sequencer.
fn genesis_stake_message(
index: usize,
sequencer_key: sequencer_stake_core::SequencerKey,
ownership_id: AccountId,
) -> Message {
let amount = system_accounts::DEFAULT_MINIMUM_SEQUENCER_STAKE;
let mover_instruction_data = lee::program::Program::serialize_instruction(
authenticated_transfer_core::Instruction::Transfer { amount },
)
.expect("Failed to serialize genesis mover instruction");
// A nonce counts how many times an account has signed. The funding account
// signed the faucet tx already, so its count starts at 1 here.
let funding_nonce = u128::try_from(index)
.expect("founding sequencer count fits in u128")
.checked_add(1)
.expect("genesis funding nonce overflow");
Message::try_new(
programs::sequencer_stake().id(),
vec![
genesis_stake_funding_account(),
ownership_id,
system_accounts::sequencer_stake_config_account_id(),
],
vec![
lee_core::account::Nonce(funding_nonce),
lee_core::account::Nonce(0),
],
sequencer_stake_core::Instruction::Stake {
sequencer_key,
amount,
mover_program_id: programs::authenticated_transfer().id(),
mover_instruction_data,
},
)
.expect("Failed to build genesis Stake message")
}
/// Signs the founding sequencer at `index`'s genesis `Stake`, for an operator
/// producing their `GenesisAction::StakeSequencer` entry.
#[must_use]
pub fn sign_genesis_stake(
index: usize,
sequencer_key: sequencer_stake_core::SequencerKey,
ownership_key: &lee::PrivateKey,
) -> lee::Signature {
let ownership_id = AccountId::from(&lee::PublicKey::new_from_private_key(ownership_key));
let message = genesis_stake_message(index, sequencer_key, ownership_id);
lee::Signature::new(ownership_key, &message.hash())
}
/// The founding sequencers' `Stake`s, funded via the faucet. Real transactions,
/// not raw state, so followers replay them instead of missing them.
fn build_stake_genesis_transactions(staked: &[FoundingStake]) -> Vec<PublicTransaction> {
if staked.is_empty() {
return Vec::new();
}
let funding_key = lee::PrivateKey::try_new(GENESIS_STAKE_FUNDING_KEY).unwrap();
let funding_public_key = lee::PublicKey::new_from_private_key(&funding_key);
let amount = system_accounts::DEFAULT_MINIMUM_SEQUENCER_STAKE;
let total = u128::try_from(staked.len())
.ok()
.and_then(|count| amount.checked_mul(count))
.expect("genesis stake total overflow");
let fund_message = Message::try_new(
programs::faucet().id(),
vec![
system_accounts::faucet_account_id(),
genesis_stake_funding_account(),
],
vec![lee_core::account::Nonce(0)],
faucet_core::Instruction::GenesisTransferDirect { amount: total },
)
.expect("Failed to build genesis funding message");
// The funding account signs even though it is only receiving. It is a brand
// new account, so the transfer claims it, and a claim needs that account's
// own signature.
let fund_witness_set =
lee::public_transaction::WitnessSet::for_message(&fund_message, &[&funding_key]);
let mut txs = vec![PublicTransaction::new(fund_message, fund_witness_set)];
for (index, (sequencer_key, ownership_public_key, signature)) in staked.iter().enumerate() {
let ownership_id = AccountId::from(ownership_public_key);
let stake_message = genesis_stake_message(index, *sequencer_key, ownership_id);
let stake_witness_set = lee::public_transaction::WitnessSet::from_raw_parts(vec![
(
lee::Signature::new(&funding_key, &stake_message.hash()),
funding_public_key.clone(),
),
(signature.clone(), ownership_public_key.clone()),
]);
// Redundant with the signature check every tx gets, but names the entry.
assert!(
stake_witness_set.is_valid_for(&stake_message),
"genesis stake signature does not match founding sequencer {index} ({})",
hex::encode(sequencer_key)
);
txs.push(PublicTransaction::new(stake_message, stake_witness_set));
}
txs
}
/// Bridge-lock holder balances configured for this zone's genesis.
fn bridge_lock_holdings(
genesis: &[GenesisAction],
) -> impl Iterator<Item = (lee::AccountId, lee::Balance)> + '_ {
genesis.iter().filter_map(|action| match action {
GenesisAction::SupplyBridgeLockHolding { holder, amount } => Some((*holder, *amount)),
GenesisAction::SupplyAccount { .. }
| GenesisAction::SupplyBridgeAccount { .. }
| GenesisAction::StakeSequencer { .. } => 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(),
}
}
fn build_bridge_deposit_tx_from_event(event: &PendingDepositEventRecord) -> Result<LeeTransaction> {
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);
// The receipt PDA carries the exactly-once check: the program reads it to
// detect a replay, so it must be in the tx's account list.
let receipt_id =
bridge_core::deposit_receipt_account_id(bridge_program_id, event.deposit_op_id.0);
let message = Message::try_new(
bridge_program_id,
vec![
system_accounts::bridge_account_id(),
recipient_vault_id,
receipt_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,
)))
}
/// Block-validity gate for a `FinalizeUnstake`, applied uniformly regardless
/// of where the transaction came from. Passes through unconditionally for
/// anything that isn't a `FinalizeUnstake` call.
fn finalize_unstake_is_includable(
state: &lee::V03State,
tx: &LeeTransaction,
live_accredited_keys: Option<&[sequencer_stake_core::SequencerKey]>,
) -> bool {
let Some(ownership_id) = finalize_unstake_ownership_account(tx) else {
return true;
};
// Without a committee snapshot there is nothing to check a FinalizeUnstake
// against, so none is includable this block.
live_accredited_keys.is_some_and(|live_accredited_keys| {
committee_discovery::finalize_unstake_is_valid(state, ownership_id, live_accredited_keys)
})
}
/// The ownership account a `FinalizeUnstake` call targets, or `None` if `tx`
/// isn't one.
fn finalize_unstake_ownership_account(tx: &LeeTransaction) -> Option<AccountId> {
let LeeTransaction::Public(tx) = tx else {
return None;
};
let message = tx.message();
if message.program_id != programs::sequencer_stake().id() {
return None;
}
match risc0_zkvm::serde::from_slice::<sequencer_stake_core::Instruction, u32>(
&message.instruction_data,
) {
Ok(sequencer_stake_core::Instruction::FinalizeUnstake) => {
message.account_ids.first().copied()
}
Ok(_) | Err(_) => None,
}
}
/// A `FinalizeUnstake` for every release `state` has pending. Whether each one
/// is actually includable is decided later, uniformly, by
/// [`finalize_unstake_is_includable`].
fn build_finalize_unstake_txs(state: &lee::V03State) -> VecDeque<LeeTransaction> {
committee_discovery::finalize_unstake_candidates(state)
.into_iter()
.filter_map(|(ownership_id, pending)| {
build_finalize_unstake_tx(ownership_id, pending)
.inspect_err(|err| warn!("Failed to build FinalizeUnstake tx: {err:#}"))
.ok()
})
.collect()
}
// Unsigned: FinalizeUnstake needs no authorization, per the program.
fn build_finalize_unstake_tx(
ownership_id: AccountId,
pending: sequencer_stake_core::PendingUnstake,
) -> Result<LeeTransaction> {
let message = Message::try_new(
programs::sequencer_stake().id(),
vec![
ownership_id,
pending.destination,
system_accounts::sequencer_stake_config_account_id(),
],
vec![],
sequencer_stake_core::Instruction::FinalizeUnstake,
)
.context("Failed to build FinalizeUnstake 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<LeeTransaction> {
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))
}
/// The cross-zone message an inbox dispatch delivers, or `None` if `tx` is not
/// a dispatch.
#[must_use]
fn extract_cross_zone_dispatch(tx: &LeeTransaction) -> Option<CrossZoneMessage> {
let LeeTransaction::Public(tx) = tx else {
return None;
};
let message = tx.message();
if message.program_id != programs::cross_zone_inbox().id() {
return None;
}
match risc0_zkvm::serde::from_slice::<cross_zone_inbox_core::Instruction, u32>(
&message.instruction_data,
) {
Ok(cross_zone_inbox_core::Instruction::Dispatch(msg)) => Some(msg),
Ok(cross_zone_inbox_core::Instruction::InitConfig(_)) | Err(_) => None,
}
}
/// The content-addressed key of the message an inbox dispatch delivers.
///
/// A delivery in an irreversible block settles its pending record, so the record
/// is dropped by identity rather than by the height it happened to land at.
#[must_use]
fn extract_cross_zone_dispatch_key(tx: &LeeTransaction) -> Option<[u8; 32]> {
extract_cross_zone_dispatch(tx).map(|msg| {
cross_zone_inbox_core::message_key(&msg.src_zone, msg.src_block_id, msg.src_tx_index)
})
}
/// The keys of the deliveries `block` carries, reporting any whose transaction
/// is not the one we recorded for that key.
///
/// The key covers `(src_zone, src_block_id, src_tx_index)` and nothing about the
/// payload, and so does the inbox's own replay check, so a sequencer that
/// publishes a dispatch with the right key and a forged payload settles our
/// correct record along with it. The forgery is caught downstream by the
/// indexer, which re-derives every delivery and halts, but the local record is
/// the last copy of what we believed and it is about to be dropped either way.
/// Saying so in the log is what makes the halt diagnosable.
fn settled_dispatch_keys(dbio: &RocksDBIO, block: &Block) -> Vec<[u8; 32]> {
let recorded = dbio.get_pending_cross_zone_dispatches().unwrap_or_default();
let (keys, forged) = classify_settled_deliveries(&recorded, block);
for key in forged {
error!(
"Cross-zone delivery {} settled with a transaction that is not the one this node recorded for that key. The message key does not cover the payload, so a peer's sequencer can publish a different delivery under it.",
hex::encode(key)
);
}
keys
}
/// Splits the deliveries `block` carries into every settled key, and the subset
/// whose transaction is not the one `recorded` holds for that key.
///
/// Separated from the logging so the detection is testable: a forged delivery
/// leaves no trace in state that differs from an honest one, precisely because
/// the key does not cover the payload.
fn classify_settled_deliveries(
recorded: &[PendingCrossZoneDispatchRecord],
block: &Block,
) -> (Vec<[u8; 32]>, Vec<[u8; 32]>) {
let mut keys = Vec::new();
let mut forged = Vec::new();
for tx in &block.body.transactions {
let Some(key) = extract_cross_zone_dispatch_key(tx) else {
continue;
};
let mismatched = recorded
.iter()
.find(|record| record.message_key == key)
.is_some_and(|record| {
borsh::to_vec(tx).is_ok_and(|encoded| encoded != record.transaction)
});
if mismatched {
forged.push(key);
}
keys.push(key);
}
(keys, forged)
}
/// Drops the records of deliveries carried by a reconstructed block.
///
/// A persist failure is only logged: the deliveries are already irreversible, so
/// the worst case is a record the next drain drops instead.
fn settle_reconstructed_deliveries(store: &SequencerStore, block: &Block) {
let keys = settled_dispatch_keys(&store.dbio(), block);
if keys.is_empty() {
return;
}
if let Err(err) = store.dbio().drop_settled_cross_zone_dispatches(&keys) {
warn!("Failed to settle reconstructed delivery records: {err:#}");
}
}
#[must_use]
fn extract_bridge_deposit_id(tx: &LeeTransaction) -> Option<HashType> {
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::<bridge_core::Instruction, u32>(&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<WithdrawArg> {
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::<bridge_core::Instruction, u32>(&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),
),
})
}
/// The reconciliation identity of one released channel note.
///
/// A `ChannelWithdrawOp` releases notes the channel already owns and carries
/// only their ids — the recipient key and value live in the note itself, which
/// neither the op nor the Bedrock Withdraw event reports. The note id is
/// therefore the one handle both sides share, and it is unique: a note is spent
/// once.
fn withdrawal_reconciliation_key(note_id: &NoteId) -> WithdrawalReconciliationKey {
let released_note_id: [u8; 32] = note_id
.as_bytes()
.as_ref()
.try_into()
.expect("`NoteId` is a 32-byte field element");
WithdrawalReconciliationKey { released_note_id }
}
/// Load key bytes from file or generate a new set if it doesn't exist.
fn load_or_create_key_bytes(path: &Path) -> Result<[u8; ED25519_SECRET_KEY_SIZE]> {
if path.exists() {
let key_bytes = std::fs::read(path)?;
key_bytes
.try_into()
.map_err(|_bytes| anyhow!("Found key with incorrect length"))
} 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(key_bytes)
}
}
/// 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<Ed25519Key> {
Ok(Ed25519Key::from_bytes(&load_or_create_key_bytes(path)?))
}
/// Load the key owning this sequencer's genesis stake, or generate one.
///
/// Only read when a solo sequencer creates the channel: a configured founding
/// set carries a signature instead, so the key never reaches the node.
pub fn load_or_create_stake_signing_key(path: &Path) -> Result<lee::PrivateKey> {
let bytes = load_or_create_key_bytes(path)?;
lee::PrivateKey::try_new(bytes).context("stake signing key file holds an invalid private key")
}
#[cfg(test)]
#[cfg(feature = "mock")]
mod tests;