use super::gas::estimate_message_gas;
use crate::lotus_json::{LotusJson, NotNullVec, lotus_json_with_self};
use crate::message::SignedMessage;
use crate::prelude::*;
use crate::rpc::error::ServerError;
use crate::rpc::types::{ApiTipsetKey, MessageSendSpec};
use crate::rpc::{ApiPaths, Ctx, Permission, RpcMethod};
use crate::shim::{
address::{Address, Protocol},
message::Message,
percent::Percent,
};
use ahash::HashSet;
use enumflags2::BitFlags;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use std::time::Duration;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "PascalCase")]
pub struct ApiMpoolConfig {
#[schemars(with = "LotusJson<Vec<Address>>")]
#[serde(with = "crate::lotus_json")]
pub priority_addrs: Vec<Address>,
pub size_limit_high: i64,
pub size_limit_low: i64,
#[serde(with = "crate::lotus_json")]
#[schemars(with = "LotusJson<Percent>")]
pub replace_by_fee_ratio: Percent,
#[schemars(with = "LotusJson<Duration>")]
#[serde(with = "crate::lotus_json")]
pub prune_cooldown: Duration,
pub gas_limit_overestimation: f64,
}
lotus_json_with_self!(ApiMpoolConfig);
pub enum MpoolGetConfig {}
impl RpcMethod<0> for MpoolGetConfig {
const NAME: &'static str = "Filecoin.MpoolGetConfig";
const PARAM_NAMES: [&'static str; 0] = [];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Read;
const DESCRIPTION: &'static str = "Returns a copy of the current mpool config.";
type Params = ();
type Ok = ApiMpoolConfig;
async fn handle(
ctx: Ctx,
(): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let cfg = ctx.mpool.config();
Ok(ApiMpoolConfig {
priority_addrs: cfg.priority_addrs,
size_limit_high: cfg.size_limit_high,
size_limit_low: cfg.size_limit_low,
replace_by_fee_ratio: cfg.replace_by_fee_ratio,
prune_cooldown: cfg.prune_cooldown,
gas_limit_overestimation: cfg.gas_limit_overestimation,
})
}
}
pub enum MpoolGetNonce {}
impl RpcMethod<1> for MpoolGetNonce {
const NAME: &'static str = "Filecoin.MpoolGetNonce";
const PARAM_NAMES: [&'static str; 1] = ["address"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Read;
const DESCRIPTION: &'static str = "Returns the current nonce for the specified address.";
type Params = (Address,);
type Ok = u64;
async fn handle(
ctx: Ctx,
(address,): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
Ok(ctx.mpool.get_sequence(&address).await?)
}
}
pub enum MpoolPending {}
impl RpcMethod<1> for MpoolPending {
const NAME: &'static str = "Filecoin.MpoolPending";
const PARAM_NAMES: [&'static str; 1] = ["tipsetKey"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Read;
const DESCRIPTION: &'static str = "Returns the pending messages for a given tipset.";
type Params = (ApiTipsetKey,);
type Ok = NotNullVec<SignedMessage>;
async fn handle(
ctx: Ctx,
(ApiTipsetKey(tipset_key),): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let mut ts = ctx
.chain_store()
.load_required_tipset_or_heaviest(&tipset_key)?;
let (mut pending, mpts) = ctx.mpool.pending();
if mpts.epoch() > ts.epoch() || mpts == ts {
return Ok(pending.into());
}
let mut have_cids: HashSet<_> = pending.iter().map(|m| m.cid()).collect();
loop {
if mpts.epoch() == ts.epoch() {
if mpts == ts {
break;
}
let have = ctx.mpool.messages_for_blocks(mpts.block_headers().iter())?;
have_cids.extend(have.iter().map(|m| m.cid()));
}
let msgs = ctx.mpool.messages_for_blocks(ts.block_headers().iter())?;
for m in msgs {
if have_cids.insert(m.cid()) {
pending.push(m);
}
}
if mpts.epoch() >= ts.epoch() {
break;
}
ts = ctx.chain_index().load_required_tipset(ts.parents())?;
}
Ok(pending.into())
}
}
pub enum MpoolSelect {}
impl RpcMethod<2> for MpoolSelect {
const NAME: &'static str = "Filecoin.MpoolSelect";
const PARAM_NAMES: [&'static str; 2] = ["tipsetKey", "ticketQuality"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Read;
const DESCRIPTION: &'static str =
"Returns a list of pending messages for inclusion in the next block.";
type Params = (ApiTipsetKey, f64);
type Ok = Vec<SignedMessage>;
async fn handle(
ctx: Ctx,
(ApiTipsetKey(tipset_key), ticket_quality): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let ts = ctx
.chain_store()
.load_required_tipset_or_heaviest(&tipset_key)?;
Ok(ctx.mpool.select_messages(&ts, ticket_quality)?)
}
}
pub enum MpoolPush {}
impl RpcMethod<1> for MpoolPush {
const NAME: &'static str = "Filecoin.MpoolPush";
const PARAM_NAMES: [&'static str; 1] = ["message"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Write;
const DESCRIPTION: &'static str = "Adds a signed message to the message pool.";
type Params = (SignedMessage,);
type Ok = Cid;
async fn handle(
ctx: Ctx,
(message,): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let cid = ctx.mpool.push(message).await?;
Ok(cid)
}
}
pub enum MpoolBatchPush {}
impl RpcMethod<1> for MpoolBatchPush {
const NAME: &'static str = "Filecoin.MpoolBatchPush";
const PARAM_NAMES: [&'static str; 1] = ["messages"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Write;
const DESCRIPTION: &'static str = "Adds a set of signed messages to the message pool.";
type Params = (Vec<SignedMessage>,);
type Ok = Vec<Cid>;
async fn handle(
ctx: Ctx,
(messages,): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let mut cids = vec![];
for msg in messages {
cids.push(ctx.mpool.push(msg).await?);
}
Ok(cids)
}
}
pub enum MpoolPushUntrusted {}
impl RpcMethod<1> for MpoolPushUntrusted {
const NAME: &'static str = "Filecoin.MpoolPushUntrusted";
const PARAM_NAMES: [&'static str; 1] = ["message"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Write;
const DESCRIPTION: &'static str =
"Adds a message to the message pool with verification checks.";
type Params = (SignedMessage,);
type Ok = Cid;
async fn handle(
ctx: Ctx,
(message,): Self::Params,
_: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let cid = ctx.mpool.push_untrusted(message).await?;
Ok(cid)
}
}
pub enum MpoolBatchPushUntrusted {}
impl RpcMethod<1> for MpoolBatchPushUntrusted {
const NAME: &'static str = "Filecoin.MpoolBatchPushUntrusted";
const PARAM_NAMES: [&'static str; 1] = ["messages"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Write;
const DESCRIPTION: &'static str =
"Adds a set of messages to the message pool with additional verification checks.";
type Params = (Vec<SignedMessage>,);
type Ok = Vec<Cid>;
async fn handle(
ctx: Ctx,
(messages,): Self::Params,
ext: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
MpoolBatchPush::handle(ctx, (messages,), ext).await
}
}
pub enum MpoolPushMessage {}
impl RpcMethod<2> for MpoolPushMessage {
const NAME: &'static str = "Filecoin.MpoolPushMessage";
const PARAM_NAMES: [&'static str; 2] = ["message", "sendSpec"];
const API_PATHS: BitFlags<ApiPaths> = ApiPaths::all();
const PERMISSION: Permission = Permission::Sign;
const DESCRIPTION: &'static str =
"Assigns a nonce, signs, and pushes a message to the mempool.";
type Params = (Message, Option<MessageSendSpec>);
type Ok = SignedMessage;
async fn handle(
ctx: Ctx,
(message, send_spec): Self::Params,
extensions: &http::Extensions,
) -> Result<Self::Ok, ServerError> {
let from = message.from;
let heaviest_tipset = ctx.chain_store().heaviest_tipset();
let key_addr = ctx
.state_manager
.resolve_to_deterministic_address(from, &heaviest_tipset)
.await?;
if message.sequence != 0 {
return Err(anyhow::anyhow!(
"Expected nonce for MpoolPushMessage is 0, and will be calculated for you"
)
.into());
}
let _sender_guard = ctx.mpool_locker.take_lock(key_addr).await;
let mut message =
estimate_message_gas(&ctx, message, send_spec, Default::default()).await?;
if message.gas_premium > message.gas_fee_cap {
return Err(anyhow::anyhow!(
"After estimation, gas premium is greater than gas fee cap"
)
.into());
}
if from.protocol() == Protocol::ID {
message.from = key_addr;
}
let balance =
super::wallet::WalletBalance::handle(ctx.clone(), (message.from,), extensions).await?;
let required_funds = &message.value + &message.gas_fee_cap * message.gas_limit;
if balance < required_funds {
return Err(anyhow::anyhow!(
"mpool push: not enough funds: {balance} < {required_funds}",
)
.into());
}
let key = crate::key_management::Key::try_from(crate::key_management::try_find(
&key_addr,
&ctx.keystore.as_ref().read(),
)?)?;
let eth_chain_id = ctx.chain_config().eth_chain_id;
let smsg = ctx
.nonce_tracker
.sign_and_push(&ctx.mpool, message, &key, eth_chain_id)
.await?;
Ok(smsg)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::blocks::{Block, CachingBlockHeader, FullTipset, RawBlockHeader, Tipset};
use crate::chain::ChainStore;
use crate::chain_sync::TipsetValidator;
use crate::rpc::RPCState;
use crate::rpc::test_utils::chain_store;
use crate::shim::crypto::{SECP_SIG_LEN, Signature};
use crate::test_utils::dummy_ticket;
use fvm_ipld_blockstore::Blockstore;
fn secp_message(sequence: u64) -> SignedMessage {
SignedMessage::new_unchecked(
Message {
from: Address::new_id(100),
to: Address::new_id(101),
sequence,
..Default::default()
},
Signature::new_secp256k1(vec![0; SECP_SIG_LEN]),
)
}
fn tipset_at(
db: &impl Blockstore,
parent: &Tipset,
epoch: i64,
ticket: u8,
bls: &[Message],
secp: &[SignedMessage],
) -> Tipset {
let fts = FullTipset::new([Block {
header: CachingBlockHeader::new(RawBlockHeader {
parents: parent.key().clone(),
epoch,
messages: TipsetValidator::compute_msg_root(db, bls, secp).unwrap(),
ticket: dummy_ticket(ticket),
..Default::default()
}),
bls_messages: bls.to_vec(),
secp_messages: secp.to_vec(),
}])
.unwrap();
fts.persist(db).unwrap();
fts.into_tipset()
}
fn tipset_on(
db: &impl Blockstore,
parent: &Tipset,
ticket: u8,
bls: &[Message],
secp: &[SignedMessage],
) -> Tipset {
tipset_at(db, parent, parent.epoch() + 1, ticket, bls, secp)
}
fn ctx_on(cs: ChainStore, mpool_ts: &Tipset) -> Arc<RPCState> {
cs.set_heaviest_tipset(mpool_ts.clone()).unwrap();
let (ctx, _) = RPCState::for_tests(cs).unwrap();
assert_eq!(
&ctx.mpool.pending().1,
mpool_ts,
"the pool must adopt the heaviest tipset"
);
ctx
}
async fn pending_at(ctx: Arc<RPCState>, ts: &Tipset) -> Vec<SignedMessage> {
let NotNullVec(pending) = MpoolPending::handle(
ctx,
(ApiTipsetKey(Some(ts.key().clone())),),
&Default::default(),
)
.await
.unwrap();
pending
}
#[tokio::test]
async fn merges_messages_of_a_same_height_fork() {
let cs = chain_store();
let genesis = cs.genesis_tipset();
let shared = secp_message(0);
let only_in_fork = secp_message(1);
let mpool_ts = tipset_on(cs.db(), &genesis, 1, &[], std::slice::from_ref(&shared));
let fork_ts = tipset_on(cs.db(), &genesis, 2, &[], &[shared, only_in_fork.clone()]);
let ctx = ctx_on(cs, &mpool_ts);
assert_eq!(pending_at(ctx, &fork_ts).await, vec![only_in_fork]);
}
#[tokio::test]
async fn walks_back_to_the_mpool_tipset() {
let cs = chain_store();
let genesis = cs.genesis_tipset();
let in_mpool_ts = secp_message(0);
let in_child = secp_message(1);
let mpool_ts = tipset_on(
cs.db(),
&genesis,
1,
&[],
std::slice::from_ref(&in_mpool_ts),
);
let child_ts = tipset_on(cs.db(), &mpool_ts, 2, &[], std::slice::from_ref(&in_child));
let ctx = ctx_on(cs, &mpool_ts);
assert_eq!(pending_at(ctx, &child_ts).await, vec![in_child]);
}
#[tokio::test]
async fn merges_across_a_null_round_past_the_mpool_tipset() {
let cs = chain_store();
let genesis = cs.genesis_tipset();
let only_in_ts = secp_message(1);
let base = tipset_on(cs.db(), &genesis, 5, &[], &[]);
let mpool_ts = tipset_on(cs.db(), &base, 1, &[], &[secp_message(0)]);
let ts = tipset_at(cs.db(), &base, 3, 3, &[], std::slice::from_ref(&only_in_ts));
let ctx = ctx_on(cs, &mpool_ts);
assert_eq!(pending_at(ctx, &ts).await, vec![only_in_ts]);
}
#[tokio::test]
async fn does_not_merge_at_or_behind_the_mpool_tipset() {
let cs = chain_store();
let genesis = cs.genesis_tipset();
let mpool_ts = tipset_on(cs.db(), &genesis, 1, &[], &[secp_message(0)]);
let ctx = ctx_on(cs, &mpool_ts);
assert!(pending_at(ctx.clone(), &mpool_ts).await.is_empty());
assert!(pending_at(ctx, &genesis).await.is_empty());
}
#[tokio::test]
async fn skips_bls_messages_with_an_uncached_signature() {
let cs = chain_store();
let genesis = cs.genesis_tipset();
let only_in_fork = secp_message(1);
let mpool_ts = tipset_on(cs.db(), &genesis, 1, &[], &[]);
let fork_ts = tipset_on(
cs.db(),
&genesis,
2,
&[Message {
from: Address::new_id(200),
to: Address::new_id(201),
..Default::default()
}],
std::slice::from_ref(&only_in_fork),
);
let ctx = ctx_on(cs, &mpool_ts);
assert_eq!(pending_at(ctx, &fork_ts).await, vec![only_in_fork]);
}
}