From f1c6aecb3d4a4ebb2cd3c9f6a58de20b019418e2 Mon Sep 17 00:00:00 2001 From: 0xfnzero <0xfnzero@users.noreply.github.com> Date: Fri, 17 Jul 2026 03:04:10 +0800 Subject: [PATCH] fix(pumpswap): preserve upgraded event fields --- Cargo.toml | 5 +- examples/pumpswap_with_metrics.rs | 29 +++--- .../event_parser/core/merger_event.rs | 99 +++++++++++++++---- .../event_parser/protocols/pumpswap/events.rs | 43 ++++++++ .../event_parser/protocols/pumpswap/types.rs | 1 + src/streaming/parser_sdk_bridge/mod.rs | 8 +- 6 files changed, 149 insertions(+), 36 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index e9b61b9..0285d17 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -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.6.0", git = "https://github.com/0xfnzero/sol-parser-sdk", rev = "1002ffeeb3a65aff0a7a38bcc834ef01a9aebf1f", 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" diff --git a/examples/pumpswap_with_metrics.rs b/examples/pumpswap_with_metrics.rs index 5f8c7cf..d0c09c6 100644 --- a/examples/pumpswap_with_metrics.rs +++ b/examples/pumpswap_with_metrics.rs @@ -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> { 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> { 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(()) } diff --git a/src/streaming/event_parser/core/merger_event.rs b/src/streaming/event_parser/core/merger_event.rs index 4a847bc..c88530c 100644 --- a/src/streaming/event_parser/core/merger_event.rs +++ b/src/streaming/event_parser/core/merger_event.rs @@ -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,16 +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.buyback_fee_basis_points = cpie.buyback_fee_basis_points; - e.buyback_fee = cpie.buyback_fee; - e.virtual_quote_reserves = cpie.virtual_quote_reserves; - e.can_boost = cpie.can_boost; - e.base_supply = cpie.base_supply; - 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); @@ -390,14 +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.buyback_fee_basis_points = cpie.buyback_fee_basis_points; - e.buyback_fee = cpie.buyback_fee; - e.virtual_quote_reserves = cpie.virtual_quote_reserves; - e.can_boost = cpie.can_boost; - e.base_supply = cpie.base_supply; - 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); @@ -947,6 +981,35 @@ mod tests { 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(); diff --git a/src/streaming/event_parser/protocols/pumpswap/events.rs b/src/streaming/event_parser/protocols/pumpswap/events.rs index f3297a4..31ca38f 100755 --- a/src/streaming/event_parser/protocols/pumpswap/events.rs +++ b/src/streaming/event_parser/protocols/pumpswap/events.rs @@ -44,10 +44,15 @@ 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, @@ -110,10 +115,15 @@ 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, @@ -251,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 { diff --git a/src/streaming/event_parser/protocols/pumpswap/types.rs b/src/streaming/event_parser/protocols/pumpswap/types.rs index c1f7f3a..7e86409 100644 --- a/src/streaming/event_parser/protocols/pumpswap/types.rs +++ b/src/streaming/event_parser/protocols/pumpswap/types.rs @@ -33,6 +33,7 @@ pub struct Pool { pub coin_creator: Pubkey, pub is_mayhem_mode: bool, pub is_cashback_coin: bool, + #[serde(default)] pub virtual_quote_reserves: i128, } diff --git a/src/streaming/parser_sdk_bridge/mod.rs b/src/streaming/parser_sdk_bridge/mod.rs index 61a5e0a..78c9ef6 100644 --- a/src/streaming/parser_sdk_bridge/mod.rs +++ b/src/streaming/parser_sdk_bridge/mod.rs @@ -245,7 +245,7 @@ mod tests { #[test] fn converts_pumpswap_buy_boost_tail() { let event = PbPumpSwapBuy { - virtual_quote_reserves: -123, + virtual_quote_reserves: i128::MIN, buyback_fee_basis_points: 11, buyback_fee: 22, can_boost: true, @@ -258,7 +258,7 @@ mod tests { let DexEvent::PumpSwapBuyEvent(converted) = converted else { panic!("expected PumpSwapBuyEvent"); }; - assert_eq!(converted.virtual_quote_reserves, -123); + 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); @@ -268,7 +268,7 @@ mod tests { #[test] fn converts_pumpswap_sell_boost_tail() { let event = PbPumpSwapSell { - virtual_quote_reserves: 123, + virtual_quote_reserves: i128::MAX, buyback_fee_basis_points: 44, buyback_fee: 55, can_boost: true, @@ -281,7 +281,7 @@ mod tests { let DexEvent::PumpSwapSellEvent(converted) = converted else { panic!("expected PumpSwapSellEvent"); }; - assert_eq!(converted.virtual_quote_reserves, 123); + 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);