mirror of
https://github.com/0xfnzero/solana-streamer.git
synced 2026-07-27 17:37:45 +00:00
Merge dev for solana-streamer-sdk 2.0.0
This commit is contained in:
+5
-2
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "solana-streamer-sdk"
|
||||
version = "1.5.16"
|
||||
version = "2.0.0"
|
||||
edition = "2021"
|
||||
authors = ["William <byteblock6@gmail.com>", "sgxiang <sgxiang@gmail.com>", "wei <1415121722@qq.com>"]
|
||||
repository = "https://github.com/0xfnzero/solana-streamer"
|
||||
@@ -21,7 +21,7 @@ sdk-perf-stats = ["sol-parser-sdk/perf-stats"]
|
||||
sdk-ultra-perf = ["sol-parser-sdk/ultra-perf"]
|
||||
|
||||
[dependencies]
|
||||
sol-parser-sdk = { version = "0.5.15", default-features = false }
|
||||
sol-parser-sdk = { version = "0.6.0", git = "https://github.com/0xfnzero/sol-parser-sdk", rev = "995d88991b56234a23fc1d0911fdd33caa063c67", default-features = false }
|
||||
solana-sdk = "3.0.0"
|
||||
solana-client = "3.1.12"
|
||||
solana-transaction-status = "3.1.12"
|
||||
@@ -50,3 +50,6 @@ spl-token-2022 = { version = "10.0.0", default-features = false, features = ["no
|
||||
spl-token-group-interface = "=0.7.1"
|
||||
solana-commitment-config = { version = "3.1.1", features = ["serde"] }
|
||||
tonic-prost = "0.14.5"
|
||||
|
||||
[dev-dependencies]
|
||||
serde_json = "1.0.149"
|
||||
|
||||
@@ -3,10 +3,14 @@
|
||||
//! Usage: cargo run --example pumpswap_with_metrics --release
|
||||
|
||||
use solana_streamer_sdk::streaming::event_parser::protocols::pumpswap::parser::PUMPSWAP_PROGRAM_ID;
|
||||
use solana_streamer_sdk::streaming::event_parser::{DexEvent, Protocol};
|
||||
use solana_streamer_sdk::streaming::event_parser::{
|
||||
common::filter::EventTypeFilter, common::EventType, DexEvent, Protocol,
|
||||
};
|
||||
use solana_streamer_sdk::streaming::grpc::ClientConfig;
|
||||
use solana_streamer_sdk::streaming::yellowstone_grpc::{
|
||||
AccountFilter, TransactionFilter, YellowstoneGrpc,
|
||||
use solana_streamer_sdk::streaming::yellowstone_grpc::{TransactionFilter, YellowstoneGrpc};
|
||||
use std::sync::{
|
||||
atomic::{AtomicU64, Ordering},
|
||||
Arc,
|
||||
};
|
||||
|
||||
#[tokio::main]
|
||||
@@ -29,22 +33,20 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
account_exclude: vec![],
|
||||
account_required: vec![],
|
||||
};
|
||||
let account_filter = AccountFilter {
|
||||
account: vec![],
|
||||
owner: vec![PUMPSWAP_PROGRAM_ID.to_string()],
|
||||
filters: vec![],
|
||||
};
|
||||
|
||||
let callback = |event: DexEvent| {
|
||||
println!("Event: {:?}", event.metadata().event_type);
|
||||
let event_count = Arc::new(AtomicU64::new(0));
|
||||
let callback_count = event_count.clone();
|
||||
let callback = move |_event: DexEvent| {
|
||||
callback_count.fetch_add(1, Ordering::Relaxed);
|
||||
};
|
||||
let event_filter =
|
||||
EventTypeFilter::include_only(vec![EventType::PumpSwapBuy, EventType::PumpSwapSell]);
|
||||
|
||||
grpc.subscribe_events_immediate(
|
||||
vec![Protocol::PumpSwap],
|
||||
None,
|
||||
vec![transaction_filter],
|
||||
vec![account_filter],
|
||||
None,
|
||||
vec![],
|
||||
Some(event_filter),
|
||||
None,
|
||||
callback,
|
||||
)
|
||||
@@ -53,5 +55,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
println!("Press Ctrl+C to stop...\n");
|
||||
tokio::signal::ctrl_c().await?;
|
||||
grpc.stop().await;
|
||||
println!("Processed {} PumpSwap trade events", event_count.load(Ordering::Relaxed));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
+1122
File diff suppressed because it is too large
Load Diff
@@ -42,6 +42,20 @@ fn fill_i64(to: &mut i64, from: i64) {
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn overwrite_u64_if_present(to: &mut u64, from: u64) {
|
||||
if from != 0 {
|
||||
*to = from;
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn overwrite_i128_if_present(to: &mut i128, from: i128) {
|
||||
if from != 0 {
|
||||
*to = from;
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn fill_string(to: &mut String, from: String) {
|
||||
if to.is_empty() && !from.is_empty() {
|
||||
@@ -338,11 +352,27 @@ pub fn merge(instruction_event: &mut DexEvent, cpi_log_event: DexEvent) {
|
||||
e.total_claimed_tokens = cpie.total_claimed_tokens;
|
||||
e.current_sol_volume = cpie.current_sol_volume;
|
||||
e.last_update_timestamp = cpie.last_update_timestamp;
|
||||
e.min_base_amount_out = cpie.min_base_amount_out;
|
||||
e.ix_name = cpie.ix_name.clone();
|
||||
e.cashback_fee_basis_points = cpie.cashback_fee_basis_points;
|
||||
e.cashback = cpie.cashback;
|
||||
e.is_pump_pool = cpie.is_pump_pool;
|
||||
overwrite_u64_if_present(&mut e.min_base_amount_out, cpie.min_base_amount_out);
|
||||
if !cpie.ix_name.is_empty() {
|
||||
e.ix_name = cpie.ix_name;
|
||||
}
|
||||
overwrite_u64_if_present(
|
||||
&mut e.cashback_fee_basis_points,
|
||||
cpie.cashback_fee_basis_points,
|
||||
);
|
||||
overwrite_u64_if_present(&mut e.cashback, cpie.cashback);
|
||||
overwrite_u64_if_present(
|
||||
&mut e.buyback_fee_basis_points,
|
||||
cpie.buyback_fee_basis_points,
|
||||
);
|
||||
overwrite_u64_if_present(&mut e.buyback_fee, cpie.buyback_fee);
|
||||
overwrite_i128_if_present(
|
||||
&mut e.virtual_quote_reserves,
|
||||
cpie.virtual_quote_reserves,
|
||||
);
|
||||
e.can_boost |= cpie.can_boost;
|
||||
overwrite_u64_if_present(&mut e.base_supply, cpie.base_supply);
|
||||
e.is_pump_pool |= cpie.is_pump_pool;
|
||||
fill_pubkey(&mut e.base_mint, cpie.base_mint);
|
||||
fill_pubkey(&mut e.quote_mint, cpie.quote_mint);
|
||||
fill_pubkey(&mut e.pool_base_token_account, cpie.pool_base_token_account);
|
||||
@@ -385,9 +415,23 @@ pub fn merge(instruction_event: &mut DexEvent, cpi_log_event: DexEvent) {
|
||||
e.coin_creator = cpie.coin_creator;
|
||||
e.coin_creator_fee_basis_points = cpie.coin_creator_fee_basis_points;
|
||||
e.coin_creator_fee = cpie.coin_creator_fee;
|
||||
e.cashback_fee_basis_points = cpie.cashback_fee_basis_points;
|
||||
e.cashback = cpie.cashback;
|
||||
e.is_pump_pool = cpie.is_pump_pool;
|
||||
overwrite_u64_if_present(
|
||||
&mut e.cashback_fee_basis_points,
|
||||
cpie.cashback_fee_basis_points,
|
||||
);
|
||||
overwrite_u64_if_present(&mut e.cashback, cpie.cashback);
|
||||
overwrite_u64_if_present(
|
||||
&mut e.buyback_fee_basis_points,
|
||||
cpie.buyback_fee_basis_points,
|
||||
);
|
||||
overwrite_u64_if_present(&mut e.buyback_fee, cpie.buyback_fee);
|
||||
overwrite_i128_if_present(
|
||||
&mut e.virtual_quote_reserves,
|
||||
cpie.virtual_quote_reserves,
|
||||
);
|
||||
e.can_boost |= cpie.can_boost;
|
||||
overwrite_u64_if_present(&mut e.base_supply, cpie.base_supply);
|
||||
e.is_pump_pool |= cpie.is_pump_pool;
|
||||
fill_pubkey(&mut e.base_mint, cpie.base_mint);
|
||||
fill_pubkey(&mut e.quote_mint, cpie.quote_mint);
|
||||
fill_pubkey(&mut e.pool_base_token_account, cpie.pool_base_token_account);
|
||||
@@ -858,6 +902,7 @@ mod tests {
|
||||
use crate::streaming::event_parser::protocols::pumpfun::events::{
|
||||
PumpFeesShareholder, PumpFunTradeEvent,
|
||||
};
|
||||
use crate::streaming::event_parser::protocols::pumpswap::events::PumpSwapSellEvent;
|
||||
|
||||
#[test]
|
||||
fn pumpfun_merge_keeps_instruction_context_and_copies_latest_trade_tail() {
|
||||
@@ -912,6 +957,59 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pumpswap_merge_preserves_signed_virtual_quote_reserves() {
|
||||
let mut instruction_event = DexEvent::PumpSwapSellEvent(PumpSwapSellEvent::default());
|
||||
let log_event = DexEvent::PumpSwapSellEvent(PumpSwapSellEvent {
|
||||
virtual_quote_reserves: -500,
|
||||
buyback_fee_basis_points: 11,
|
||||
buyback_fee: 22,
|
||||
can_boost: true,
|
||||
base_supply: 33,
|
||||
..Default::default()
|
||||
});
|
||||
|
||||
merge(&mut instruction_event, log_event);
|
||||
|
||||
let DexEvent::PumpSwapSellEvent(event) = instruction_event else {
|
||||
panic!("expected PumpSwapSellEvent");
|
||||
};
|
||||
assert_eq!(event.virtual_quote_reserves, -500);
|
||||
assert_eq!(event.buyback_fee_basis_points, 11);
|
||||
assert_eq!(event.buyback_fee, 22);
|
||||
assert!(event.can_boost);
|
||||
assert_eq!(event.base_supply, 33);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pumpswap_merge_does_not_erase_present_tail_with_legacy_defaults() {
|
||||
let mut instruction_event = DexEvent::PumpSwapSellEvent(PumpSwapSellEvent {
|
||||
cashback_fee_basis_points: 10,
|
||||
cashback: 20,
|
||||
buyback_fee_basis_points: 30,
|
||||
buyback_fee: 40,
|
||||
virtual_quote_reserves: -500,
|
||||
can_boost: true,
|
||||
base_supply: 60,
|
||||
is_pump_pool: true,
|
||||
..Default::default()
|
||||
});
|
||||
|
||||
merge(&mut instruction_event, DexEvent::PumpSwapSellEvent(PumpSwapSellEvent::default()));
|
||||
|
||||
let DexEvent::PumpSwapSellEvent(event) = instruction_event else {
|
||||
panic!("expected PumpSwapSellEvent");
|
||||
};
|
||||
assert_eq!(event.cashback_fee_basis_points, 10);
|
||||
assert_eq!(event.cashback, 20);
|
||||
assert_eq!(event.buyback_fee_basis_points, 30);
|
||||
assert_eq!(event.buyback_fee, 40);
|
||||
assert_eq!(event.virtual_quote_reserves, -500);
|
||||
assert!(event.can_boost);
|
||||
assert_eq!(event.base_supply, 60);
|
||||
assert!(event.is_pump_pool);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pumpfun_merge_replaces_sol_quote_sentinel_with_real_quote_mint() {
|
||||
let quote_mint = Pubkey::new_unique();
|
||||
|
||||
@@ -44,6 +44,16 @@ pub struct PumpSwapBuyEvent {
|
||||
pub ix_name: String,
|
||||
pub cashback_fee_basis_points: u64,
|
||||
pub cashback: u64,
|
||||
#[serde(default)]
|
||||
pub buyback_fee_basis_points: u64,
|
||||
#[serde(default)]
|
||||
pub buyback_fee: u64,
|
||||
#[serde(default)]
|
||||
pub virtual_quote_reserves: i128,
|
||||
#[serde(default)]
|
||||
pub can_boost: bool,
|
||||
#[serde(default)]
|
||||
pub base_supply: u64,
|
||||
#[borsh(skip)]
|
||||
pub is_pump_pool: bool,
|
||||
#[borsh(skip)]
|
||||
@@ -105,6 +115,16 @@ pub struct PumpSwapSellEvent {
|
||||
pub coin_creator_fee: u64,
|
||||
pub cashback_fee_basis_points: u64,
|
||||
pub cashback: u64,
|
||||
#[serde(default)]
|
||||
pub buyback_fee_basis_points: u64,
|
||||
#[serde(default)]
|
||||
pub buyback_fee: u64,
|
||||
#[serde(default)]
|
||||
pub virtual_quote_reserves: i128,
|
||||
#[serde(default)]
|
||||
pub can_boost: bool,
|
||||
#[serde(default)]
|
||||
pub base_supply: u64,
|
||||
#[borsh(skip)]
|
||||
pub is_pump_pool: bool,
|
||||
#[borsh(skip)]
|
||||
@@ -241,6 +261,39 @@ pub struct PumpSwapWithdrawEvent {
|
||||
|
||||
pub const PUMP_SWAP_WITHDRAW_EVENT_LOG_SIZE: usize = 248;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn pumpswap_event_serde_preserves_signed_i128_extremes() {
|
||||
for value in [i128::MIN, -1, i128::MAX] {
|
||||
let event = PumpSwapBuyEvent { virtual_quote_reserves: value, ..Default::default() };
|
||||
let json = serde_json::to_string(&event).unwrap();
|
||||
let decoded: PumpSwapBuyEvent = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(decoded.virtual_quote_reserves, value);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pumpswap_event_serde_defaults_new_upgrade_tail() {
|
||||
let mut value = serde_json::to_value(PumpSwapSellEvent::default()).unwrap();
|
||||
let object = value.as_object_mut().unwrap();
|
||||
object.remove("buyback_fee_basis_points");
|
||||
object.remove("buyback_fee");
|
||||
object.remove("virtual_quote_reserves");
|
||||
object.remove("can_boost");
|
||||
object.remove("base_supply");
|
||||
|
||||
let decoded: PumpSwapSellEvent = serde_json::from_value(value).unwrap();
|
||||
assert_eq!(decoded.buyback_fee_basis_points, 0);
|
||||
assert_eq!(decoded.buyback_fee, 0);
|
||||
assert_eq!(decoded.virtual_quote_reserves, 0);
|
||||
assert!(!decoded.can_boost);
|
||||
assert_eq!(decoded.base_supply, 0);
|
||||
}
|
||||
}
|
||||
|
||||
/// 全局配置
|
||||
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize, BorshDeserialize)]
|
||||
pub struct PumpSwapGlobalConfigAccountEvent {
|
||||
|
||||
@@ -33,13 +33,13 @@ pub struct Pool {
|
||||
pub coin_creator: Pubkey,
|
||||
pub is_mayhem_mode: bool,
|
||||
pub is_cashback_coin: bool,
|
||||
/// On-chain reserved tail (7 bytes); keep in sync with pump_amm pool account layout.
|
||||
pub reserved: [u8; 7],
|
||||
#[serde(default)]
|
||||
pub virtual_quote_reserves: i128,
|
||||
}
|
||||
|
||||
/// Legacy pool account body (before `is_cashback_coin` + reserved).
|
||||
pub const POOL_BODY_LEGACY: usize = 1 + 2 + 32 * 6 + 8 + 32 + 1;
|
||||
/// Current pool account body including flags and reserved.
|
||||
pub const POOL_BODY: usize = POOL_BODY_LEGACY + 1 + 7;
|
||||
/// Legacy allocated pool account body.
|
||||
pub const POOL_BODY_LEGACY: usize = 244;
|
||||
/// Current serialized pool account body including signed virtual quote reserves.
|
||||
pub const POOL_BODY: usize = 253;
|
||||
|
||||
pub const POOL_SIZE: usize = POOL_BODY;
|
||||
|
||||
@@ -185,7 +185,7 @@ pub(crate) fn pumpswap_pool_from_pb(p: sol_parser_sdk::core::events::PumpSwapPoo
|
||||
coin_creator: p.coin_creator,
|
||||
is_mayhem_mode: p.is_mayhem_mode,
|
||||
is_cashback_coin: p.is_cashback_coin,
|
||||
reserved: [0u8; 7],
|
||||
virtual_quote_reserves: p.virtual_quote_reserves,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -43,10 +43,10 @@ mod tests {
|
||||
use sol_parser_sdk::core::events::{
|
||||
EventMetadata, MeteoraDlmmSwapEvent as PbDlmmSwap, OrcaWhirlpoolSwapEvent as PbOrcaSwap,
|
||||
PumpFunCreateV2TokenEvent as PbPumpCreateV2, PumpFunTradeEvent as PbPumpTrade,
|
||||
PumpSwapCreatePoolEvent as PbPumpSwapCreatePool, PumpSwapPool as PbPumpSwapPool,
|
||||
PumpSwapPoolAccountEvent as PbPumpSwapPoolAccount,
|
||||
RaydiumLaunchlabTradeEvent as PbBonkTrade, TokenInfoEvent as PbTokenInfo,
|
||||
TradeDirection as PbBonkDir,
|
||||
PumpSwapBuyEvent as PbPumpSwapBuy, PumpSwapCreatePoolEvent as PbPumpSwapCreatePool,
|
||||
PumpSwapPool as PbPumpSwapPool, PumpSwapPoolAccountEvent as PbPumpSwapPoolAccount,
|
||||
PumpSwapSellEvent as PbPumpSwapSell, RaydiumLaunchlabTradeEvent as PbBonkTrade,
|
||||
TokenInfoEvent as PbTokenInfo, TradeDirection as PbBonkDir,
|
||||
};
|
||||
use sol_parser_sdk::DexEvent as PbDexEvent;
|
||||
use solana_sdk::{pubkey::Pubkey, signature::Signature};
|
||||
@@ -242,6 +242,52 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn converts_pumpswap_buy_boost_tail() {
|
||||
let event = PbPumpSwapBuy {
|
||||
virtual_quote_reserves: i128::MIN,
|
||||
buyback_fee_basis_points: 11,
|
||||
buyback_fee: 22,
|
||||
can_boost: true,
|
||||
base_supply: 33,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let converted =
|
||||
convert_parser_event(PbDexEvent::PumpSwapBuy(event), None, 999).expect("convert");
|
||||
let DexEvent::PumpSwapBuyEvent(converted) = converted else {
|
||||
panic!("expected PumpSwapBuyEvent");
|
||||
};
|
||||
assert_eq!(converted.virtual_quote_reserves, i128::MIN);
|
||||
assert_eq!(converted.buyback_fee_basis_points, 11);
|
||||
assert_eq!(converted.buyback_fee, 22);
|
||||
assert!(converted.can_boost);
|
||||
assert_eq!(converted.base_supply, 33);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn converts_pumpswap_sell_boost_tail() {
|
||||
let event = PbPumpSwapSell {
|
||||
virtual_quote_reserves: i128::MAX,
|
||||
buyback_fee_basis_points: 44,
|
||||
buyback_fee: 55,
|
||||
can_boost: true,
|
||||
base_supply: 66,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let converted =
|
||||
convert_parser_event(PbDexEvent::PumpSwapSell(event), None, 999).expect("convert");
|
||||
let DexEvent::PumpSwapSellEvent(converted) = converted else {
|
||||
panic!("expected PumpSwapSellEvent");
|
||||
};
|
||||
assert_eq!(converted.virtual_quote_reserves, i128::MAX);
|
||||
assert_eq!(converted.buyback_fee_basis_points, 44);
|
||||
assert_eq!(converted.buyback_fee, 55);
|
||||
assert!(converted.can_boost);
|
||||
assert_eq!(converted.base_supply, 66);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn converts_pumpswap_pool_account_cashback_flag() {
|
||||
let pool_pubkey = Pubkey::new_unique();
|
||||
@@ -272,6 +318,7 @@ mod tests {
|
||||
coin_creator: Pubkey::new_unique(),
|
||||
is_mayhem_mode: true,
|
||||
is_cashback_coin: true,
|
||||
virtual_quote_reserves: -777,
|
||||
},
|
||||
};
|
||||
|
||||
@@ -283,6 +330,7 @@ mod tests {
|
||||
assert_eq!(st.pubkey, pool_pubkey);
|
||||
assert!(st.pool.is_mayhem_mode);
|
||||
assert!(st.pool.is_cashback_coin);
|
||||
assert_eq!(st.pool.virtual_quote_reserves, -777);
|
||||
}
|
||||
_ => panic!("expected PumpSwapPoolAccountEvent"),
|
||||
}
|
||||
|
||||
@@ -523,6 +523,11 @@ pub(crate) fn pumpswap_buy_full_from_parser(
|
||||
ix_name: b.ix_name,
|
||||
cashback_fee_basis_points: b.cashback_fee_basis_points,
|
||||
cashback: b.cashback,
|
||||
buyback_fee_basis_points: b.buyback_fee_basis_points,
|
||||
buyback_fee: b.buyback_fee,
|
||||
virtual_quote_reserves: b.virtual_quote_reserves,
|
||||
can_boost: b.can_boost,
|
||||
base_supply: b.base_supply,
|
||||
is_pump_pool: b.is_pump_pool,
|
||||
base_mint: b.base_mint,
|
||||
quote_mint: b.quote_mint,
|
||||
@@ -569,6 +574,11 @@ pub(crate) fn pumpswap_sell_full_from_parser(
|
||||
coin_creator_fee: s.coin_creator_fee,
|
||||
cashback_fee_basis_points: s.cashback_fee_basis_points,
|
||||
cashback: s.cashback,
|
||||
buyback_fee_basis_points: s.buyback_fee_basis_points,
|
||||
buyback_fee: s.buyback_fee,
|
||||
virtual_quote_reserves: s.virtual_quote_reserves,
|
||||
can_boost: s.can_boost,
|
||||
base_supply: s.base_supply,
|
||||
is_pump_pool: s.is_pump_pool,
|
||||
base_mint: s.base_mint,
|
||||
quote_mint: s.quote_mint,
|
||||
|
||||
Reference in New Issue
Block a user