feat(sequencer): add RPC Server Actor

This commit is contained in:
Daniil Polyakov
2026-08-12 22:34:14 +03:00
parent 436908e08f
commit 4ae9b30736
14 changed files with 236 additions and 141 deletions
Generated
+28 -13
View File
@@ -2298,7 +2298,7 @@ dependencies = [
"clap",
"json-pretty-compact",
"sequencer_core_metrics",
"sequencer_service_metrics",
"sequencer_rpc_server_actor_metrics",
"serde",
"serde_json",
]
@@ -9577,38 +9577,53 @@ dependencies = [
]
[[package]]
name = "sequencer_service"
name = "sequencer_rpc_server_actor"
version = "0.1.0"
dependencies = [
"anyhow",
"borsh",
"bytesize",
"common",
"jsonrpsee",
"kameo",
"lee",
"log",
"programs",
"sequencer_core",
"sequencer_executor_actor",
"sequencer_rpc_server_actor_metrics",
"sequencer_service_protocol",
"sequencer_service_rpc",
]
[[package]]
name = "sequencer_rpc_server_actor_metrics"
version = "0.1.0"
dependencies = [
"metrics",
]
[[package]]
name = "sequencer_service"
version = "0.1.0"
dependencies = [
"anyhow",
"bytesize",
"clap",
"common",
"env_logger",
"futures",
"hex",
"jsonrpsee",
"lee",
"log",
"mempool",
"metrics-exporter-prometheus",
"programs",
"sequencer_core",
"sequencer_service_metrics",
"sequencer_service_protocol",
"sequencer_service_rpc",
"tokio",
"tokio-util",
]
[[package]]
name = "sequencer_service_metrics"
version = "0.1.0"
dependencies = [
"metrics",
]
[[package]]
name = "sequencer_service_protocol"
version = "0.1.0"
+4 -2
View File
@@ -22,8 +22,9 @@ members = [
"lez/sequencer/service",
"lez/sequencer/service/protocol",
"lez/sequencer/service/rpc",
"lez/sequencer/service/metrics",
"lez/sequencer/actors/executor",
"lez/sequencer/actors/rpc_server",
"lez/sequencer/actors/rpc_server/metrics",
"lez/indexer/core",
"lez/indexer/service",
"lez/indexer/service/protocol",
@@ -85,8 +86,9 @@ sequencer_core = { path = "lez/sequencer/core" }
sequencer_core_metrics = { path = "lez/sequencer/core/metrics" }
sequencer_service_protocol = { path = "lez/sequencer/service/protocol" }
sequencer_service_rpc = { path = "lez/sequencer/service/rpc" }
sequencer_service_metrics = { path = "lez/sequencer/service/metrics" }
sequencer_executor_actor = { path = "lez/sequencer/actors/executor" }
sequencer_rpc_server_actor = { path = "lez/sequencer/actors/rpc_server" }
sequencer_rpc_server_actor_metrics = { path = "lez/sequencer/actors/rpc_server/metrics" }
sequencer_service = { path = "lez/sequencer/service" }
indexer_core = { path = "lez/indexer/core" }
indexer_service = { path = "lez/indexer/service" }
@@ -0,0 +1,25 @@
[package]
name = "sequencer_rpc_server_actor"
version = "0.1.0"
edition = "2024"
license = { workspace = true }
[lints]
workspace = true
[dependencies]
lee.workspace = true
common.workspace = true
programs.workspace = true
sequencer_core.workspace = true
sequencer_service_protocol.workspace = true
sequencer_service_rpc = { workspace = true, features = ["server"] }
sequencer_rpc_server_actor_metrics = { workspace = true, features = ["record"] }
sequencer_executor_actor.workspace = true
kameo.workspace = true
log.workspace = true
jsonrpsee.workspace = true
borsh.workspace = true
anyhow.workspace = true
bytesize.workspace = true
@@ -1,5 +1,5 @@
[package]
name = "sequencer_service_metrics"
name = "sequencer_rpc_server_actor_metrics"
version = "0.1.0"
edition = "2024"
license = { workspace = true }
@@ -1,4 +1,4 @@
//! This crate provides all metrics exposed by the sequencer service crate.
//! This crate provides all metrics exposed by RPC Server Actor.
#[cfg(feature = "record")]
pub use record::*;
+103
View File
@@ -0,0 +1,103 @@
//! RPC Server Actor server RPC queries and retranslates them to Executor.
use std::net::SocketAddr;
use anyhow::{Context as _, Result};
use bytesize::ByteSize;
use jsonrpsee::{server::ServerHandle, tracing::warn};
use kameo::{Actor, actor::ActorRef, mailbox::Signal, message::Message};
use log::info;
use sequencer_service_rpc::RpcServer as _;
use crate::protocol::{GetAddress, GetAddressReply};
pub mod protocol;
mod service;
const REQUEST_BODY_MAX_SIZE: ByteSize = ByteSize::mib(10);
pub struct RpcServerActor {
server_handle: Option<ServerHandle>,
addr: SocketAddr,
}
impl RpcServerActor {
pub async fn new(
executor_ref: ActorRef<sequencer_executor_actor::ExecutorActor>,
listen_addr: SocketAddr,
max_block_size: u64,
) -> Result<Self> {
let server = jsonrpsee::server::ServerBuilder::with_config(
jsonrpsee::server::ServerConfigBuilder::new()
.max_request_body_size(
u32::try_from(REQUEST_BODY_MAX_SIZE.as_u64())
.expect("REQUEST_BODY_MAX_SIZE should be less than u32::MAX"),
)
.build(),
)
.build(listen_addr)
.await
.context("Failed to build RPC server")?;
let addr = server
.local_addr()
.context("Failed to get local address of RPC server")?;
info!("Starting RPC Server on {addr}");
let service = service::Service::new(executor_ref, max_block_size);
let server_handle = server.start(service.into_rpc());
Ok(Self {
server_handle: Some(server_handle),
addr,
})
}
}
impl Actor for RpcServerActor {
type Args = Self;
type Error = anyhow::Error;
async fn on_start(args: Self::Args, _actor_ref: ActorRef<Self>) -> Result<Self, Self::Error> {
Ok(args)
}
async fn next(
&mut self,
_actor_ref: kameo::prelude::WeakActorRef<Self>,
_mailbox_rx: &mut kameo::prelude::MailboxReceiver<Self>,
) -> Result<Option<Signal<Self>>, Self::Error> {
if let Some(server_handle) = self.server_handle.take() {
server_handle.stopped().await;
warn!("RPC server stopped");
}
Ok(Some(Signal::Stop))
}
async fn on_stop(
&mut self,
_actor_ref: kameo::prelude::WeakActorRef<Self>,
_reason: kameo::prelude::ActorStopReason,
) -> Result<(), Self::Error> {
if let Some(server_handle) = self.server_handle.take() {
server_handle.stop()?;
info!("RPC server stopped");
}
Ok(())
}
}
impl Message<GetAddress> for RpcServerActor {
type Reply = GetAddressReply;
async fn handle(
&mut self,
GetAddress: GetAddress,
_ctx: &mut kameo::prelude::Context<Self, Self::Reply>,
) -> Self::Reply {
GetAddressReply { addr: self.addr }
}
}
@@ -0,0 +1,10 @@
use std::net::SocketAddr;
use kameo::Reply;
pub struct GetAddress;
#[derive(Reply)]
pub struct GetAddressReply {
pub addr: SocketAddr,
}
@@ -1,50 +1,40 @@
use std::{collections::BTreeMap, sync::Arc};
use std::collections::BTreeMap;
use common::transaction::LeeTransaction;
use jsonrpsee::{
core::async_trait,
types::{ErrorCode, ErrorObjectOwned},
};
use lee;
use kameo::actor::ActorRef;
use log::{error, warn};
use mempool::MemPoolHandle;
use sequencer_core::{
DbError, SequencerCore, TransactionOrigin, block_publisher::BlockPublisherTrait,
};
use sequencer_service_protocol::{
Account, AccountId, Block, BlockId, ChannelId, Commitment, CommitmentSetDigest,
CrossZoneDeadLetter, CrossZoneDeadLetterReport, HashType, MembershipProof, Nonce, ProgramId,
};
use tokio::sync::Mutex;
const NOT_FOUND_ERROR_CODE: i32 = -31999;
pub struct SequencerService<BC: BlockPublisherTrait> {
sequencer: Arc<Mutex<SequencerCore<BC>>>,
mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>,
pub struct Service {
executor_ref: ActorRef<sequencer_executor_actor::ExecutorActor>,
max_block_size: u64,
}
impl<BC: BlockPublisherTrait> SequencerService<BC> {
pub const fn new(
sequencer: Arc<Mutex<SequencerCore<BC>>>,
mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>,
impl Service {
pub fn new(
executor_ref: ActorRef<sequencer_executor_actor::ExecutorActor>,
max_block_size: u64,
) -> Self {
sequencer_rpc_server_actor_metrics::init();
Self {
sequencer,
mempool_handle,
executor_ref,
max_block_size,
}
}
}
#[async_trait]
impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::RpcServer
for SequencerService<BC>
{
impl sequencer_service_rpc::RpcServer for Service {
async fn send_transaction(&self, tx: LeeTransaction) -> Result<HashType, ErrorObjectOwned> {
sequencer_service_metrics::increment_submitted_transactions_total();
sequencer_rpc_server_actor_metrics::increment_submitted_transactions_total();
let tx_hash = tx.hash();
@@ -97,14 +87,17 @@ impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::Rpc
};
let authenticated_tx = res.await.inspect_err(|err| {
sequencer_service_metrics::increment_before_mempool_failed_transactions_total();
sequencer_rpc_server_actor_metrics::increment_before_mempool_failed_transactions_total(
);
error!("Transaction failed before reaching mempool: {err:#?}");
})?;
self.mempool_handle
.push((TransactionOrigin::User, authenticated_tx))
self.executor_ref
.tell(sequencer_executor_actor::protocol::Transaction {
transaction: authenticated_tx,
})
.await
.expect("Mempool is closed, this is a bug");
.map_err(internal_error)?;
Ok(tx_hash)
}
@@ -114,11 +107,10 @@ impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::Rpc
}
async fn get_block(&self, block_id: BlockId) -> Result<Option<Block>, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
sequencer
.block_store()
.get_block_at_id(block_id)
.map_err(|err| internal_error(&err))
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetBlock { block_id })
.await
.map_err(internal_error)
}
async fn get_block_range(
@@ -126,74 +118,64 @@ impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::Rpc
start_block_id: BlockId,
end_block_id: BlockId,
) -> Result<Vec<Block>, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
(start_block_id..=end_block_id)
.map(|block_id| {
let block = sequencer
.block_store()
.get_block_at_id(block_id)
.map_err(|err| internal_error(&err))?;
block.ok_or_else(|| {
ErrorObjectOwned::owned(
NOT_FOUND_ERROR_CODE,
format!("Block with id {block_id} not found"),
None::<()>,
)
})
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetBlockRange {
range: (start_block_id..=end_block_id),
})
.collect::<Result<Vec<_>, _>>()
.await
.map_err(internal_error)
}
async fn get_last_block_id(&self) -> Result<BlockId, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
Ok(sequencer.chain_height())
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetLastBlockId)
.await
.map_err(internal_error)
}
async fn get_account_balance(&self, account_id: AccountId) -> Result<u128, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
let balance = sequencer.with_state(|state| state.get_account_by_id(account_id).balance);
Ok(balance)
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetAccountBalance { account_id })
.await
.map_err(internal_error)
}
async fn get_transaction(
&self,
tx_hash: HashType,
) -> Result<Option<(LeeTransaction, BlockId)>, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
Ok(sequencer.block_store().get_transaction_by_hash(tx_hash))
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetTransaction { tx_hash })
.await
.map_err(internal_error)
}
async fn get_accounts_nonces(
&self,
account_ids: Vec<AccountId>,
) -> Result<Vec<Nonce>, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
let nonces = sequencer.with_state(|state| {
account_ids
.into_iter()
.map(|account_id| state.get_account_by_id(account_id).nonce)
.collect()
});
Ok(nonces)
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetAccountNonces { account_ids })
.await
.map_err(internal_error)
}
async fn get_proofs_and_root(
&self,
commitments: Vec<Commitment>,
) -> Result<(Vec<Option<MembershipProof>>, CommitmentSetDigest), ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
Ok(sequencer.with_state(|state| {
let proofs = commitments
.iter()
.map(|commitment| state.get_proof_for_commitment(commitment))
.collect();
(proofs, state.commitment_root())
}))
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetProofsAndRoot { commitments })
.await
.map_err(internal_error)
}
async fn get_account(&self, account_id: AccountId) -> Result<Account, ErrorObjectOwned> {
let sequencer = self.sequencer.lock().await;
Ok(sequencer.with_state(|state| state.get_account_by_id(account_id)))
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetAccount { account_id })
.await
.map(|reply| reply.account)
.map_err(internal_error)
}
async fn get_program_ids(&self) -> Result<BTreeMap<String, ProgramId>, ErrorObjectOwned> {
@@ -214,8 +196,11 @@ impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::Rpc
}
async fn get_channel_id(&self) -> Result<ChannelId, ErrorObjectOwned> {
let channel_id = self.sequencer.lock().await.block_publisher().channel_id();
Ok(ChannelId(*channel_id.as_ref()))
self.executor_ref
.ask(sequencer_executor_actor::protocol::GetChannelId)
.await
.map(|reply| ChannelId(reply.channel_id))
.map_err(internal_error)
}
async fn get_cross_zone_dead_letters(
@@ -245,6 +230,6 @@ impl<BC: BlockPublisherTrait + Send + Sync + 'static> sequencer_service_rpc::Rpc
}
}
fn internal_error(err: &DbError) -> ErrorObjectOwned {
fn internal_error(err: impl std::fmt::Display) -> ErrorObjectOwned {
ErrorObjectOwned::owned(ErrorCode::InternalError.code(), err.to_string(), None::<()>)
}
-5
View File
@@ -10,13 +10,9 @@ workspace = true
[dependencies]
common.workspace = true
lee.workspace = true
mempool.workspace = true
sequencer_core = { workspace = true, features = ["testnet"] }
sequencer_service_protocol.workspace = true
sequencer_service_rpc = { workspace = true, features = ["server"] }
sequencer_service_metrics = { workspace = true, features = ["record"] }
programs.workspace = true
clap = { workspace = true, features = ["derive", "env"] }
anyhow.workspace = true
@@ -29,7 +25,6 @@ tokio-util.workspace = true
jsonrpsee.workspace = true
futures.workspace = true
bytesize.workspace = true
borsh.workspace = true
[features]
default = []
-40
View File
@@ -1,30 +1,21 @@
use std::{net::SocketAddr, sync::Arc, time::Duration};
use anyhow::{Context as _, Result, anyhow};
use bytesize::ByteSize;
use common::transaction::LeeTransaction;
use futures::never::Never;
use jsonrpsee::server::ServerHandle;
use log::{error, info, warn};
use mempool::MemPoolHandle;
#[cfg(not(feature = "standalone"))]
use sequencer_core::SequencerCore;
#[cfg(feature = "standalone")]
use sequencer_core::SequencerCoreWithMockClients as SequencerCore;
pub use sequencer_core::config::*;
use sequencer_core::{
TransactionOrigin,
block_publisher::BlockPublisherTrait as _,
task_group::{StoreRelease, TaskGroup},
};
use sequencer_service_rpc::RpcServer as _;
use tokio::{sync::Mutex, task::JoinHandle};
use tokio_util::sync::CancellationToken;
pub mod service;
const REQUEST_BODY_MAX_SIZE: ByteSize = ByteSize::mib(10);
/// Handle to manage the sequencer and its tasks.
///
/// Implements `Drop` to ensure all tasks are aborted and the RPC server is stopped when dropped.
@@ -211,8 +202,6 @@ async fn wait_for_store_release(store: &StoreRelease) {
}
pub async fn run(config: SequencerConfig, listen_addr: SocketAddr) -> Result<SequencerHandle> {
sequencer_service_metrics::init();
let block_timeout = config.block_create_timeout;
let max_block_size = config.max_block_size;
@@ -254,35 +243,6 @@ pub async fn run(config: SequencerConfig, listen_addr: SocketAddr) -> Result<Seq
))
}
async fn run_server(
sequencer: Arc<Mutex<SequencerCore>>,
mempool_handle: MemPoolHandle<(TransactionOrigin, LeeTransaction)>,
listen_addr: SocketAddr,
max_block_size: u64,
) -> Result<(ServerHandle, SocketAddr)> {
let server = jsonrpsee::server::ServerBuilder::with_config(
jsonrpsee::server::ServerConfigBuilder::new()
.max_request_body_size(
u32::try_from(REQUEST_BODY_MAX_SIZE.as_u64())
.expect("REQUEST_BODY_MAX_SIZE should be less than u32::MAX"),
)
.build(),
)
.build(listen_addr)
.await
.context("Failed to build RPC server")?;
let addr = server
.local_addr()
.context("Failed to get local address of RPC server")?;
info!("Starting Sequencer Service RPC server on {addr}");
let service = service::SequencerService::new(sequencer, mempool_handle, max_block_size);
let handle = server.start(service.into_rpc());
Ok((handle, addr))
}
async fn main_loop(seq_core: Arc<Mutex<SequencerCore>>, block_timeout: Duration) -> Result<Never> {
loop {
tokio::time::sleep(block_timeout).await;
+1 -1
View File
@@ -9,7 +9,7 @@ workspace = true
[dependencies]
sequencer_core_metrics.workspace = true
sequencer_service_metrics.workspace = true
sequencer_rpc_server_actor_metrics.workspace = true
clap = { workspace = true, features = ["derive"] }
serde = { workspace = true, features = ["derive", "alloc"] }
@@ -156,9 +156,9 @@ pub fn dashboard() -> Dashboard {
// `clamp_min` keeps an idle window (nothing submitted)
// reading as 0% instead of a division by zero.
"100 * (increase({before_mempool}[$__range]) + increase({in_mempool}[$__range])) / clamp_min(increase({submitted}[$__range]), 1)",
before_mempool = sequencer_service_metrics::names::BEFORE_MEMPOOL_FAILED_TRANSACTIONS_TOTAL,
before_mempool = sequencer_rpc_server_actor_metrics::names::BEFORE_MEMPOOL_FAILED_TRANSACTIONS_TOTAL,
in_mempool = sequencer_core_metrics::names::MEMPOOL_FAILED_TRANSACTIONS_TOTAL,
submitted = sequencer_service_metrics::names::SUBMITTED_TRANSACTIONS_TOTAL,
submitted = sequencer_rpc_server_actor_metrics::names::SUBMITTED_TRANSACTIONS_TOTAL,
))
.legend("failed"),
),
@@ -169,11 +169,11 @@ pub fn dashboard() -> Dashboard {
.fill_opacity(35)
.gradient_mode(GradientMode::Opacity)
.target(rate_per_min(
sequencer_service_metrics::names::SUBMITTED_TRANSACTIONS_TOTAL,
sequencer_rpc_server_actor_metrics::names::SUBMITTED_TRANSACTIONS_TOTAL,
"submitted",
))
.target(rate_per_min(
sequencer_service_metrics::names::BEFORE_MEMPOOL_FAILED_TRANSACTIONS_TOTAL,
sequencer_rpc_server_actor_metrics::names::BEFORE_MEMPOOL_FAILED_TRANSACTIONS_TOTAL,
"failed · before mempool",
))
.target(rate_per_min(