diff --git a/Cargo.toml b/Cargo.toml index e9fa18d..71f0acd 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "solana-streamer-sdk" -version = "0.4.5" +version = "0.4.6" edition = "2021" authors = ["William ", "sgxiang ", "wei <1415121722@qq.com>"] repository = "https://github.com/0xfnzero/solana-streamer" @@ -69,6 +69,7 @@ crossbeam-queue = "0.3.12" parking_lot = "0.12.1" wide = "0.7" spl-token = "8.0.0" +spl-token-2022 = "9.0.0" [dev-dependencies] criterion = { version = "0.5", features = ["html_reports"] } \ No newline at end of file diff --git a/README.md b/README.md index 50ccc53..eff7063 100755 --- a/README.md +++ b/README.md @@ -104,14 +104,14 @@ Add the dependency to your `Cargo.toml`: ```toml # Add to your Cargo.toml -solana-streamer-sdk = { path = "./solana-streamer", version = "0.4.5" } +solana-streamer-sdk = { path = "./solana-streamer", version = "0.4.6" } ``` ### Use crates.io ```toml # Add to your Cargo.toml -solana-streamer-sdk = "0.4.5" +solana-streamer-sdk = "0.4.6" ``` ## ⚙️ Configuration System diff --git a/README_CN.md b/README_CN.md index d5cde7c..7223c80 100644 --- a/README_CN.md +++ b/README_CN.md @@ -104,14 +104,14 @@ git clone https://github.com/0xfnzero/solana-streamer ```toml # 添加到您的 Cargo.toml -solana-streamer-sdk = { path = "./solana-streamer", version = "0.4.5" } +solana-streamer-sdk = { path = "./solana-streamer", version = "0.4.6" } ``` ### 使用 crates.io ```toml # 添加到您的 Cargo.toml -solana-streamer-sdk = "0.4.5" +solana-streamer-sdk = "0.4.6" ``` ## ⚙️ 配置系统 diff --git a/examples/dynamic_subscription.rs b/examples/dynamic_subscription.rs index f5ffe5c..d6e84f9 100644 --- a/examples/dynamic_subscription.rs +++ b/examples/dynamic_subscription.rs @@ -1,11 +1,13 @@ use anyhow::Result; -use solana_streamer_sdk::streaming::yellowstone_grpc::{AccountFilter, TransactionFilter, YellowstoneGrpc}; -use solana_streamer_sdk::streaming::event_parser::Protocol; +use solana_sdk::signature::{Keypair, Signer}; use solana_streamer_sdk::streaming::event_parser::common::filter::EventTypeFilter; use solana_streamer_sdk::streaming::event_parser::common::types::EventType; -use solana_sdk::signature::{Keypair, Signer}; -use std::sync::{Arc, Mutex}; +use solana_streamer_sdk::streaming::event_parser::Protocol; +use solana_streamer_sdk::streaming::yellowstone_grpc::{ + AccountFilter, TransactionFilter, YellowstoneGrpc, +}; use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use tokio::time::sleep; @@ -20,27 +22,28 @@ const MONITORING_DURATION_SECS: u64 = 10; #[tokio::main] async fn main() -> Result<()> { env_logger::init(); - + println!("Connecting to Yellowstone gRPC at {}", GRPC_ENDPOINT); - let client = Arc::new(YellowstoneGrpc::new( - GRPC_ENDPOINT.to_string(), - API_KEY.map(|s| s.to_string()) - )?); + let client = + Arc::new(YellowstoneGrpc::new(GRPC_ENDPOINT.to_string(), API_KEY.map(|s| s.to_string()))?); let event_counter = Arc::new(AtomicU64::new(0)); let counter = event_counter.clone(); - - let callback = move |event: Box| { - let count = counter.fetch_add(1, Ordering::Relaxed); - - let protocol = match event.event_type() { - EventType::PumpFunBuy | EventType::PumpFunSell => "PumpFun", - EventType::RaydiumCpmmSwapBaseInput | EventType::RaydiumCpmmSwapBaseOutput => "RaydiumCpmm", - _ => "Unknown" + + let callback = + move |event: Box| { + let count = counter.fetch_add(1, Ordering::Relaxed); + + let protocol = match event.event_type() { + EventType::PumpFunBuy | EventType::PumpFunSell => "PumpFun", + EventType::RaydiumCpmmSwapBaseInput | EventType::RaydiumCpmmSwapBaseOutput => { + "RaydiumCpmm" + } + _ => "Unknown", + }; + + println!("Event #{}: {:11} - {:.8}...", count + 1, protocol, event.signature()); }; - - println!("Event #{}: {:11} - {:.8}...", count + 1, protocol, event.signature()); - }; println!("\n=== Phase 1: PumpFun only ==="); let pumpfun_filter = TransactionFilter { @@ -48,11 +51,8 @@ async fn main() -> Result<()> { account_exclude: vec![], account_required: vec![], }; - - let account_filter = AccountFilter { - account: vec![], - owner: vec![], - }; + + let account_filter = AccountFilter { account: vec![], owner: vec![], filters: vec![] }; let trade_event_filter = EventTypeFilter { include: vec![ EventType::PumpFunBuy, @@ -62,46 +62,52 @@ async fn main() -> Result<()> { ], }; - if let Err(e) = client.subscribe_events_immediate( - vec![Protocol::PumpFun, Protocol::RaydiumCpmm], - None, - pumpfun_filter, - account_filter, - Some(trade_event_filter), - None, - callback, - ).await { + if let Err(e) = client + .subscribe_events_immediate( + vec![Protocol::PumpFun, Protocol::RaydiumCpmm], + None, + vec![pumpfun_filter], + vec![account_filter], + Some(trade_event_filter), + None, + callback, + ) + .await + { println!("Failed to create subscription: {}", e); return Ok(()); } - println!("Subscribed to PumpFun transactions with trade event filters, monitoring for {}s...", MONITORING_DURATION_SECS); + println!( + "Subscribed to PumpFun transactions with trade event filters, monitoring for {}s...", + MONITORING_DURATION_SECS + ); sleep(Duration::from_secs(MONITORING_DURATION_SECS)).await; let phase1_count = event_counter.load(Ordering::Relaxed); println!("Phase 1: {} events", phase1_count); println!("\n=== Phase 2: PumpFun + RaydiumCpmm ==="); let multi_protocol_filter = TransactionFilter { - account_include: vec![ - PUMPFUN_PROGRAM_ID.to_string(), - RAYDIUM_CPMM_PROGRAM_ID.to_string(), - ], + account_include: vec![PUMPFUN_PROGRAM_ID.to_string(), RAYDIUM_CPMM_PROGRAM_ID.to_string()], account_exclude: vec![], account_required: vec![], }; - if let Err(e) = client.update_subscription( - multi_protocol_filter, - AccountFilter { - account: vec![], - owner: vec![], - }, - ).await { + if let Err(e) = client + .update_subscription( + vec![multi_protocol_filter], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + ) + .await + { println!("Failed to update subscription: {}", e); return Ok(()); } - println!("Updated to PumpFun + RaydiumCpmm transactions, monitoring for {}s...", MONITORING_DURATION_SECS); + println!( + "Updated to PumpFun + RaydiumCpmm transactions, monitoring for {}s...", + MONITORING_DURATION_SECS + ); sleep(Duration::from_secs(MONITORING_DURATION_SECS)).await; let phase2_count = event_counter.load(Ordering::Relaxed); println!("Phase 2: {} events", phase2_count - phase1_count); @@ -113,19 +119,22 @@ async fn main() -> Result<()> { account_required: vec![], }; - if let Err(e) = client.update_subscription( - raydium_cpmm_filter, - AccountFilter { - account: vec![], - owner: vec![], - }, - ).await { + if let Err(e) = client + .update_subscription( + vec![raydium_cpmm_filter], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + ) + .await + { println!("Failed to update subscription: {}", e); return Ok(()); } sleep(Duration::from_secs(MONITORING_DURATION_SECS)).await; - println!("Updated to RaydiumCpmm transactions only, monitoring for {}s...", MONITORING_DURATION_SECS); + println!( + "Updated to RaydiumCpmm transactions only, monitoring for {}s...", + MONITORING_DURATION_SECS + ); let phase3_count = event_counter.load(Ordering::Relaxed); println!("Phase 3: {} events", phase3_count - phase2_count); @@ -136,19 +145,22 @@ async fn main() -> Result<()> { account_required: vec![], }; - if let Err(e) = client.update_subscription( - pumpfun_only_filter, - AccountFilter { - account: vec![], - owner: vec![], - }, - ).await { + if let Err(e) = client + .update_subscription( + vec![pumpfun_only_filter], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + ) + .await + { println!("Failed to update subscription: {}", e); return Ok(()); } sleep(Duration::from_secs(MONITORING_DURATION_SECS)).await; - println!("Updated to PumpFun transactions only, monitoring for {}s...", MONITORING_DURATION_SECS); + println!( + "Updated to PumpFun transactions only, monitoring for {}s...", + MONITORING_DURATION_SECS + ); let phase4_count = event_counter.load(Ordering::Relaxed); println!("Phase 4: {} events", phase4_count - phase3_count); @@ -158,46 +170,46 @@ async fn main() -> Result<()> { account_exclude: vec![], account_required: vec![], }; - - if let Err(e) = client.update_subscription( - empty_filter, - AccountFilter { - account: vec![], - owner: vec![], - }, - ).await { + + if let Err(e) = client + .update_subscription( + vec![empty_filter], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + ) + .await + { println!("Failed to update subscription: {}", e); return Ok(()); } - + sleep(Duration::from_secs(MONITORING_DURATION_SECS)).await; - println!("Updated to all transactions (no filters), monitoring for {}s...", MONITORING_DURATION_SECS); + println!( + "Updated to all transactions (no filters), monitoring for {}s...", + MONITORING_DURATION_SECS + ); let phase5_count = event_counter.load(Ordering::Relaxed); println!("Phase 5: {} events", phase5_count - phase4_count); println!("\n=== Phase 6: Silence ==="); - + let random_keypair_1 = Keypair::new(); let random_keypair_2 = Keypair::new(); let random_pubkey_1 = random_keypair_1.pubkey(); let random_pubkey_2 = random_keypair_2.pubkey(); - + let silence_filter = TransactionFilter { account_include: vec![], account_exclude: vec![], - account_required: vec![ - random_pubkey_1.to_string(), - random_pubkey_2.to_string(), - ], + account_required: vec![random_pubkey_1.to_string(), random_pubkey_2.to_string()], }; - if let Err(e) = client.update_subscription( - silence_filter, - AccountFilter { - account: vec![], - owner: vec![], - }, - ).await { + if let Err(e) = client + .update_subscription( + vec![silence_filter], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + ) + .await + { println!("Failed to update subscription: {}", e); return Ok(()); } @@ -222,44 +234,46 @@ async fn main() -> Result<()> { let final_count = event_counter.load(Ordering::Relaxed); let events_during_silence = final_count - before_silence; - + if events_during_silence == 0 { println!("Phase 6: 0 events (immediate filter application)"); } else if let Ok(last_time) = last_event_time.lock() { let propagation_time = last_time.duration_since(start_time); - println!("Phase 6: {} events during propagation, filter took {}ms", - events_during_silence, propagation_time.as_millis()); + println!( + "Phase 6: {} events during propagation, filter took {}ms", + events_during_silence, + propagation_time.as_millis() + ); } println!("\n=== Phase 7: Shutdown ==="); - - let shutdown_client = Arc::new(YellowstoneGrpc::new( - GRPC_ENDPOINT.to_string(), - API_KEY.map(|s| s.to_string()) - )?); - + + let shutdown_client = + Arc::new(YellowstoneGrpc::new(GRPC_ENDPOINT.to_string(), API_KEY.map(|s| s.to_string()))?); + let shutdown_event_counter = Arc::new(AtomicU64::new(0)); let shutdown_counter = shutdown_event_counter.clone(); - let shutdown_callback = move |_event: Box| { - shutdown_counter.fetch_add(1, Ordering::Relaxed); - }; + let shutdown_callback = + move |_event: Box| { + shutdown_counter.fetch_add(1, Ordering::Relaxed); + }; - if let Err(e) = shutdown_client.subscribe_events_immediate( - vec![Protocol::PumpFun, Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - shutdown_callback, - ).await { + if let Err(e) = shutdown_client + .subscribe_events_immediate( + vec![Protocol::PumpFun, Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + shutdown_callback, + ) + .await + { println!("Failed to subscribe shutdown client: {}", e); return Ok(()); } @@ -278,7 +292,7 @@ async fn main() -> Result<()> { if during_stop > 0 { println!(" {} events received during stop()", during_stop); } - + let last_event_time = Arc::new(Mutex::new(stop_time)); let last_event_time_clone = last_event_time.clone(); let mut last_count = post_stop_count; @@ -296,167 +310,169 @@ async fn main() -> Result<()> { let final_count = shutdown_event_counter.load(Ordering::Relaxed); let after_stop = final_count - post_stop_count; - + if after_stop == 0 { println!("Phase 7: Clean shutdown - no events after stop()"); } else if let Ok(last_time) = last_event_time.lock() { let post_stop_duration = last_time.duration_since(stop_time); let silence_duration = Instant::now().duration_since(*last_time); - println!("Phase 7: {} events arrived up to {}ms after stop(), then silent for {}ms", - after_stop, post_stop_duration.as_millis(), silence_duration.as_millis()); + println!( + "Phase 7: {} events arrived up to {}ms after stop(), then silent for {}ms", + after_stop, + post_stop_duration.as_millis(), + silence_duration.as_millis() + ); } println!("\n=== Subscription enforcement ==="); - let test_callback = |_event: Box| {}; + let test_callback = + |_event: Box| {}; - match client.subscribe_events_immediate( - vec![Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - test_callback, - ).await { + match client + .subscribe_events_immediate( + vec![Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + test_callback, + ) + .await + { Ok(_) => println!("ERROR: Same client created second subscription"), Err(e) if e.to_string().contains("Already subscribed") => { println!("✓ Single subscription enforcement working"); - }, + } Err(e) => println!("Unexpected error: {}", e), } - let client2 = Arc::new(YellowstoneGrpc::new( - GRPC_ENDPOINT.to_string(), - API_KEY.map(|s| s.to_string()) - )?); + let client2 = + Arc::new(YellowstoneGrpc::new(GRPC_ENDPOINT.to_string(), API_KEY.map(|s| s.to_string()))?); let client2_counter = Arc::new(AtomicU64::new(0)); let counter2 = client2_counter.clone(); - let client2_callback = move |_event: Box| { - counter2.fetch_add(1, Ordering::Relaxed); - }; + let client2_callback = + move |_event: Box| { + counter2.fetch_add(1, Ordering::Relaxed); + }; - match client2.subscribe_events_immediate( - vec![Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - client2_callback, - ).await { + match client2 + .subscribe_events_immediate( + vec![Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + client2_callback, + ) + .await + { Ok(_) => { sleep(Duration::from_millis(500)).await; let count = client2_counter.load(Ordering::Relaxed); println!("✓ Second client: {} events", count); client2.stop().await; - }, + } Err(e) => println!("ERROR: Second client failed: {}", e), } println!("\n=== Advanced subscription enforcement ==="); - - let test_callback_advanced = |_event: Box| {}; - let client3 = Arc::new(YellowstoneGrpc::new( - GRPC_ENDPOINT.to_string(), - API_KEY.map(|s| s.to_string()) - )?); + let test_callback_advanced = + |_event: Box| {}; + + let client3 = + Arc::new(YellowstoneGrpc::new(GRPC_ENDPOINT.to_string(), API_KEY.map(|s| s.to_string()))?); // First subscription should succeed - match client3.subscribe_events_immediate( - vec![Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - test_callback_advanced, - ).await { + match client3 + .subscribe_events_immediate( + vec![Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + test_callback_advanced, + ) + .await + { Ok(_) => { // Second subscription attempt on same client should fail - match client3.subscribe_events_immediate( - vec![Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - |_| {}, - ).await { + match client3 + .subscribe_events_immediate( + vec![Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + |_| {}, + ) + .await + { Ok(_) => println!("ERROR: Same client created second advanced subscription"), Err(e) if e.to_string().contains("Already subscribed") => { println!("✓ Advanced single subscription enforcement working"); - }, + } Err(e) => println!("Unexpected error: {}", e), } - }, + } Err(e) => println!("ERROR: First advanced subscription failed: {}", e), } // Test that a second client can subscribe using advanced method - let client4 = Arc::new(YellowstoneGrpc::new( - GRPC_ENDPOINT.to_string(), - API_KEY.map(|s| s.to_string()) - )?); + let client4 = + Arc::new(YellowstoneGrpc::new(GRPC_ENDPOINT.to_string(), API_KEY.map(|s| s.to_string()))?); let client4_counter = Arc::new(AtomicU64::new(0)); let counter4 = client4_counter.clone(); - let client4_callback = move |_event: Box| { - counter4.fetch_add(1, Ordering::Relaxed); - }; + let client4_callback = + move |_event: Box| { + counter4.fetch_add(1, Ordering::Relaxed); + }; - match client4.subscribe_events_immediate( - vec![Protocol::RaydiumCpmm], - None, - TransactionFilter { - account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], - account_exclude: vec![], - account_required: vec![], - }, - AccountFilter { - account: vec![], - owner: vec![], - }, - None, - None, - client4_callback, - ).await { + match client4 + .subscribe_events_immediate( + vec![Protocol::RaydiumCpmm], + None, + vec![TransactionFilter { + account_include: vec![RAYDIUM_CPMM_PROGRAM_ID.to_string()], + account_exclude: vec![], + account_required: vec![], + }], + vec![AccountFilter { account: vec![], owner: vec![], filters: vec![] }], + None, + None, + client4_callback, + ) + .await + { Ok(_) => { sleep(Duration::from_millis(500)).await; let count = client4_counter.load(Ordering::Relaxed); println!("✓ Second client (advanced): {} events", count); client4.stop().await; - }, + } Err(e) => println!("ERROR: Second client (advanced) failed: {}", e), } diff --git a/examples/grpc_example.rs b/examples/grpc_example.rs index ae0a466..b904628 100644 --- a/examples/grpc_example.rs +++ b/examples/grpc_example.rs @@ -3,7 +3,7 @@ use solana_streamer_sdk::{ streaming::{ event_parser::{ common::EventType, - core::account_event_parser::{TokenInfoEvent, NonceAccountEvent, TokenAccountEvent}, + core::account_event_parser::{NonceAccountEvent, TokenAccountEvent, TokenInfoEvent}, protocols::{ bonk::{ parser::BONK_PROGRAM_ID, BonkGlobalConfigAccountEvent, BonkMigrateToAmmEvent, @@ -105,7 +105,7 @@ async fn test_grpc() -> Result<(), Box> { }; // Listen to account data belonging to owner programs -> account event monitoring - let account_filter = AccountFilter { account: vec![], owner: account_include.clone() }; + let account_filter = AccountFilter { account: vec![], owner: account_include.clone(), filters: vec![] }; // Event filtering // No event filtering, includes all events @@ -121,8 +121,8 @@ async fn test_grpc() -> Result<(), Box> { grpc.subscribe_events_immediate( protocols, None, - transaction_filter, - account_filter, + vec![transaction_filter], + vec![account_filter], event_type_filter, None, callback, diff --git a/examples/nonce_listen_example.rs b/examples/nonce_listen_example.rs index bb3321f..fc4c52f 100644 --- a/examples/nonce_listen_example.rs +++ b/examples/nonce_listen_example.rs @@ -42,7 +42,8 @@ async fn test_grpc() -> Result<(), Box> { let nonce_account = "use_your_nonce_account_here".to_string(); // Listen to account data belonging to owner programs -> account event monitoring - let account_filter = AccountFilter { account: vec![nonce_account], owner: vec![] }; + let account_filter = + AccountFilter { account: vec![nonce_account], owner: vec![], filters: vec![] }; // Event filtering let event_type_filter = Some(EventTypeFilter { include: vec![EventType::NonceAccount] }); @@ -53,8 +54,8 @@ async fn test_grpc() -> Result<(), Box> { grpc.subscribe_events_immediate( protocols, None, - transaction_filter, - account_filter, + vec![transaction_filter], + vec![account_filter], event_type_filter, None, callback, diff --git a/examples/pumpswap_pool_account_listen_example.rs b/examples/pumpswap_pool_account_listen_example.rs new file mode 100644 index 0000000..b7fe92e --- /dev/null +++ b/examples/pumpswap_pool_account_listen_example.rs @@ -0,0 +1,119 @@ +use std::str::FromStr; + +use solana_sdk::pubkey::Pubkey; +use solana_streamer_sdk::{ + match_event, + streaming::{ + event_parser::{ + common::{filter::EventTypeFilter, EventType}, + core::account_event_parser::TokenAccountEvent, + UnifiedEvent, + }, + grpc::ClientConfig, + yellowstone_grpc::{AccountFilter, TransactionFilter}, + YellowstoneGrpc, + }, +}; +use yellowstone_grpc_proto::geyser::{ + subscribe_request_filter_accounts_filter::Filter, + subscribe_request_filter_accounts_filter_memcmp::Data, SubscribeRequestFilterAccountsFilter, + SubscribeRequestFilterAccountsFilterMemcmp, +}; + +#[tokio::main] +async fn main() -> Result<(), Box> { + println!("Starting Yellowstone gRPC Streamer..."); + test_grpc().await?; + Ok(()) +} + +async fn test_grpc() -> Result<(), Box> { + println!("Subscribing to Yellowstone gRPC events..."); + // Create low-latency configuration + let mut config: ClientConfig = ClientConfig::low_latency(); + // Enable performance monitoring, has performance overhead, disabled by default + config.enable_metrics = true; + let grpc = YellowstoneGrpc::new_with_config( + "https://solana-yellowstone-grpc.publicnode.com:443".to_string(), + None, + config, + )?; + println!("GRPC client created successfully"); + let callback = create_event_callback(); + // Will try to parse corresponding protocol events from transactions + let protocols = vec![]; + println!("Protocols to monitor: {:?}", protocols); + // Filter accounts + let account_include = vec![]; + let account_exclude = vec![]; + let account_required = vec![]; + + // Listen to transaction data + let transaction_filter = + TransactionFilter { account_include, account_exclude, account_required }; + + // Pump.fun AMM (PUMP-USDC) Market + let pump_usdc = Pubkey::from_str("2uF4Xh61rDwxnG9woyxsVQP7zuA6kLFpb3NvnRQeoiSd").unwrap(); + let wsol_deepseekai = Pubkey::from_str("BJAjivuMVANjpRWtrRfcxzGhnMSywBN19Sa4jAzWxXDx").unwrap(); + + // Listen to account data belonging to owner programs -> account event monitoring + let pump_usdc_account_filter = AccountFilter { + account: vec![], + owner: vec![], + filters: vec![SubscribeRequestFilterAccountsFilter { + filter: Some(Filter::Memcmp(SubscribeRequestFilterAccountsFilterMemcmp { + offset: 32, + data: Some(Data::Bytes(pump_usdc.to_bytes().to_vec())), + })), + }], + }; + let wsol_deepseekai_account_filter = AccountFilter { + account: vec![], + owner: vec![], + filters: vec![SubscribeRequestFilterAccountsFilter { + filter: Some(Filter::Memcmp(SubscribeRequestFilterAccountsFilterMemcmp { + offset: 32, + data: Some(Data::Bytes(wsol_deepseekai.to_bytes().to_vec())), + })), + }], + }; + + // Event filtering + let event_type_filter = Some(EventTypeFilter { include: vec![EventType::TokenAccount] }); + + println!("Starting to listen for events, press Ctrl+C to stop..."); + println!("Starting subscription..."); + + grpc.subscribe_events_immediate( + protocols.clone(), + None, + vec![transaction_filter.clone()], + vec![pump_usdc_account_filter.clone(), wsol_deepseekai_account_filter.clone()], + event_type_filter.clone(), + None, + callback, + ) + .await?; + + // 支持 stop 方法,测试代码 - 异步1000秒之后停止 + let grpc_clone = grpc.clone(); + tokio::spawn(async move { + tokio::time::sleep(std::time::Duration::from_secs(1000)).await; + grpc_clone.stop().await; + }); + + println!("Waiting for Ctrl+C to stop..."); + tokio::signal::ctrl_c().await?; + + Ok(()) +} + +fn create_event_callback() -> impl Fn(Box) { + |event: Box| { + match_event!(event, { + TokenAccountEvent => |e: TokenAccountEvent| { + println!("TokenAccount: {:?} amount: {:?}", e.pubkey, e.amount); + }, + }); + } +} diff --git a/examples/token_balance_listen_example.rs b/examples/token_balance_listen_example.rs index ee47c5e..ac1cb54 100644 --- a/examples/token_balance_listen_example.rs +++ b/examples/token_balance_listen_example.rs @@ -43,7 +43,8 @@ async fn test_grpc() -> Result<(), Box> { let account_to_listen = "use_your_token_account_here".to_string(); // Listen to account data belonging to owner programs -> account event monitoring - let account_filter = AccountFilter { account: vec![account_to_listen], owner: vec![] }; + let account_filter = + AccountFilter { account: vec![account_to_listen], owner: vec![], filters: vec![] }; // Event filtering let event_type_filter = Some(EventTypeFilter { include: vec![EventType::TokenAccount] }); @@ -54,8 +55,8 @@ async fn test_grpc() -> Result<(), Box> { grpc.subscribe_events_immediate( protocols.clone(), None, - transaction_filter.clone(), - account_filter.clone(), + vec![transaction_filter.clone()], + vec![account_filter.clone()], event_type_filter.clone(), None, callback, diff --git a/examples/token_decimals_listen_example.rs b/examples/token_decimals_listen_example.rs index 38bc9e3..1abde5d 100644 --- a/examples/token_decimals_listen_example.rs +++ b/examples/token_decimals_listen_example.rs @@ -47,7 +47,8 @@ async fn test_grpc() -> Result<(), Box> { let account_to_listen = "EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v".to_string(); // Listen to account data belonging to owner programs -> account event monitoring - let account_filter = AccountFilter { account: vec![account_to_listen], owner: vec![] }; + let account_filter = + AccountFilter { account: vec![account_to_listen], owner: vec![], filters: vec![] }; // Event filtering let event_type_filter = Some(EventTypeFilter { include: vec![EventType::TokenAccount] }); @@ -58,8 +59,8 @@ async fn test_grpc() -> Result<(), Box> { grpc.subscribe_events_immediate( protocols.clone(), None, - transaction_filter.clone(), - account_filter.clone(), + vec![transaction_filter.clone()], + vec![account_filter.clone()], event_type_filter.clone(), None, callback, diff --git a/src/streaming/event_parser/core/account_event_parser.rs b/src/streaming/event_parser/core/account_event_parser.rs index 62d30c4..8713059 100644 --- a/src/streaming/event_parser/core/account_event_parser.rs +++ b/src/streaming/event_parser/core/account_event_parser.rs @@ -1,12 +1,3 @@ -use std::collections::HashMap; -use std::sync::OnceLock; - -use serde::{Deserialize, Serialize}; -use solana_account_decoder::parse_nonce::parse_nonce; -use solana_sdk::program_pack::Pack; -use solana_sdk::pubkey::Pubkey; -use spl_token::state::Account; - use crate::impl_unified_event; use crate::streaming::common::SimdUtils; use crate::streaming::event_parser::common::filter::EventTypeFilter; @@ -20,8 +11,17 @@ use crate::streaming::event_parser::protocols::raydium_clmm::parser::RAYDIUM_CLM use crate::streaming::event_parser::protocols::raydium_cpmm::parser::RAYDIUM_CPMM_PROGRAM_ID; use crate::streaming::event_parser::Protocol; use crate::streaming::grpc::AccountPretty; - -use spl_token::state::Mint; +use serde::{Deserialize, Serialize}; +use solana_account_decoder::parse_nonce::parse_nonce; +use solana_sdk::program_pack::Pack; +use solana_sdk::pubkey::Pubkey; +use spl_token::state::{Account, Mint}; +use spl_token_2022::{ + extension::StateWithExtensions, + state::{Account as Account2022, Mint as Mint2022}, +}; +use std::collections::HashMap; +use std::sync::OnceLock; /// 通用事件解析器配置 #[derive(Debug, Clone)] @@ -295,6 +295,7 @@ impl AccountEventParser { let lamports = account.lamports; let owner = account.owner; let rent_epoch = account.rent_epoch; + // Spl Token Mint if account.data.len() >= Mint::LEN { if let Ok(mint) = Mint::unpack_from_slice(&account.data) { let mut event = TokenInfoEvent { @@ -312,7 +313,32 @@ impl AccountEventParser { return Some(Box::new(event)); } } - let amount = Account::unpack(&account.data).ok().map(|info| info.amount); + // Spl Token2022 Mint + if account.data.len() >= Account2022::LEN { + if let Ok(mint) = StateWithExtensions::::unpack(&account.data) { + let mut event = TokenInfoEvent { + metadata, + pubkey, + executable, + lamports, + owner, + rent_epoch, + supply: mint.base.supply, + decimals: mint.base.decimals, + }; + let recv_delta = elapsed_micros_since(account.recv_us); + event.set_handle_us(recv_delta); + return Some(Box::new(event)); + } + } + let amount = if account.owner == spl_token_2022::ID { + StateWithExtensions::::unpack(&account.data) + .ok() + .map(|info| info.base.amount) + } else { + Account::unpack(&account.data).ok().map(|info| info.amount) + }; + let mut event = TokenAccountEvent { metadata, pubkey, diff --git a/src/streaming/grpc/subscription.rs b/src/streaming/grpc/subscription.rs index 06d4d59..9c50340 100644 --- a/src/streaming/grpc/subscription.rs +++ b/src/streaming/grpc/subscription.rs @@ -13,6 +13,8 @@ use super::types::TransactionsFilterMap; use crate::common::AnyResult; use crate::streaming::common::StreamClientConfig as ClientConfig; use crate::streaming::event_parser::common::filter::EventTypeFilter; +use crate::streaming::yellowstone_grpc::AccountFilter; +use crate::streaming::yellowstone_grpc::TransactionFilter; /// Subscription manager #[derive(Clone)] @@ -51,15 +53,14 @@ impl SubscriptionManager { impl Stream>, SubscribeRequest, )> { - let blocks_meta = if event_type_filter.is_some() - && event_type_filter.unwrap().include_block_event() - { - hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} } - } else if event_type_filter.is_none() { - hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} } - } else { - hashmap! {} - }; + let blocks_meta = + if event_type_filter.is_some() && event_type_filter.unwrap().include_block_event() { + hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} } + } else if event_type_filter.is_none() { + hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} } + } else { + hashmap! {} + }; let subscribe_request = SubscribeRequest { accounts: accounts.unwrap_or_default(), transactions: transactions.unwrap_or_default(), @@ -79,56 +80,53 @@ impl SubscriptionManager { /// Create account subscription request and return stream pub fn subscribe_with_account_request( &self, - account: Vec, - owner: Vec, + account_filter: Vec, event_type_filter: Option<&EventTypeFilter>, ) -> Option { - if account.len() == 0 && owner.len() == 0 { + if event_type_filter.is_some() && !event_type_filter.unwrap().include_account_event() { return None; } - if event_type_filter.is_some() - && !event_type_filter.unwrap().include_account_event() - { + if account_filter.len() == 0 { return None; } let mut accounts = HashMap::new(); - accounts.insert( - "".to_owned(), - SubscribeRequestFilterAccounts { - account: account, - owner: owner, - filters: vec![], - nonempty_txn_signature: None, - }, - ); + for af in account_filter { + accounts.insert( + "".to_owned(), + SubscribeRequestFilterAccounts { + account: af.account, + owner: af.owner, + filters: af.filters, + nonempty_txn_signature: None, + }, + ); + } Some(accounts) } /// Generate subscription request filter pub fn get_subscribe_request_filter( &self, - account_include: Vec, - account_exclude: Vec, - account_required: Vec, + transaction_filter: Vec, event_type_filter: Option<&EventTypeFilter>, ) -> Option { - if event_type_filter.is_some() - && !event_type_filter.unwrap().include_transaction_event() - { + if event_type_filter.is_some() && !event_type_filter.unwrap().include_transaction_event() { return None; } let mut transactions = HashMap::new(); - transactions.insert( - "client".to_string(), - SubscribeRequestFilterTransactions { - vote: Some(false), - failed: Some(false), - signature: None, - account_include, - account_exclude, - account_required, - }, - ); + for tf in transaction_filter { + transactions.insert( + "client".to_string(), + SubscribeRequestFilterTransactions { + vote: Some(false), + failed: Some(false), + signature: None, + account_include: tf.account_include, + account_exclude: tf.account_exclude, + account_required: tf.account_required, + }, + ); + } Some(transactions) } diff --git a/src/streaming/yellowstone_grpc.rs b/src/streaming/yellowstone_grpc.rs index bf30b31..12ac830 100644 --- a/src/streaming/yellowstone_grpc.rs +++ b/src/streaming/yellowstone_grpc.rs @@ -4,10 +4,8 @@ use crate::streaming::common::{ }; use crate::streaming::event_parser::common::filter::EventTypeFilter; use crate::streaming::event_parser::{Protocol, UnifiedEvent}; -use crate::streaming::grpc::{ - EventPretty, SubscriptionManager, -}; use crate::streaming::grpc::pool::factory; +use crate::streaming::grpc::{EventPretty, SubscriptionManager}; use anyhow::anyhow; use chrono::Local; use futures::channel::mpsc; @@ -18,7 +16,9 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, RwLock}; use tokio::sync::Mutex; use yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof; -use yellowstone_grpc_proto::geyser::{CommitmentLevel, SubscribeRequest, SubscribeRequestPing}; +use yellowstone_grpc_proto::geyser::{ + CommitmentLevel, SubscribeRequest, SubscribeRequestFilterAccountsFilter, SubscribeRequestPing, +}; /// 交易过滤器 #[derive(Debug, Clone)] @@ -33,6 +33,7 @@ pub struct TransactionFilter { pub struct AccountFilter { pub account: Vec, pub owner: Vec, + pub filters: Vec, } pub struct YellowstoneGrpc { @@ -48,6 +49,8 @@ pub struct YellowstoneGrpc { pub active_subscription: Arc, pub control_tx: Arc>>>, pub current_request: Arc>>, + + pub event_type_filter: Arc>>, } impl YellowstoneGrpc { @@ -86,6 +89,7 @@ impl YellowstoneGrpc { active_subscription: Arc::new(AtomicBool::new(false)), control_tx: Arc::new(tokio::sync::Mutex::new(None)), current_request: Arc::new(tokio::sync::RwLock::new(None)), + event_type_filter: Arc::new(tokio::sync::RwLock::new(None)), }) } @@ -107,7 +111,6 @@ impl YellowstoneGrpc { Self::new_with_config(endpoint, x_token, StreamClientConfig::low_latency()) } - /// 获取配置 pub fn get_config(&self) -> &StreamClientConfig { &self.config @@ -161,8 +164,8 @@ impl YellowstoneGrpc { &self, protocols: Vec, bot_wallet: Option, - transaction_filter: TransactionFilter, - account_filter: AccountFilter, + transaction_filter: Vec, + account_filter: Vec, event_type_filter: Option, commitment: Option, callback: F, @@ -170,6 +173,7 @@ impl YellowstoneGrpc { where F: Fn(Box) + Send + Sync + 'static, { + *self.event_type_filter.write().await = event_type_filter.clone(); if self .active_subscription .compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed) @@ -184,17 +188,12 @@ impl YellowstoneGrpc { metrics_handle = self.metrics_manager.start_auto_monitoring().await; } - let transactions = self.subscription_manager.get_subscribe_request_filter( - transaction_filter.account_include, - transaction_filter.account_exclude, - transaction_filter.account_required, - event_type_filter.as_ref(), - ); - let accounts = self.subscription_manager.subscribe_with_account_request( - account_filter.account, - account_filter.owner, - event_type_filter.as_ref(), - ); + let transactions = self + .subscription_manager + .get_subscribe_request_filter(transaction_filter, event_type_filter.as_ref()); + let accounts = self + .subscription_manager + .subscribe_with_account_request(account_filter, event_type_filter.as_ref()); // 订阅事件 let (mut subscribe_tx, mut stream, subscribe_request) = self @@ -322,8 +321,8 @@ impl YellowstoneGrpc { /// Returns `AnyResult<()>` on success, error on failure pub async fn update_subscription( &self, - transaction_filter: TransactionFilter, - account_filter: AccountFilter, + transaction_filter: Vec, + account_filter: Vec, ) -> AnyResult<()> { let mut control_sender = { let control_guard = self.control_tx.lock().await; @@ -349,16 +348,17 @@ impl YellowstoneGrpc { request.transactions = self .subscription_manager .get_subscribe_request_filter( - transaction_filter.account_include, - transaction_filter.account_exclude, - transaction_filter.account_required, - None, + transaction_filter, + self.event_type_filter.read().await.as_ref(), ) .unwrap_or_default(); request.accounts = self .subscription_manager - .subscribe_with_account_request(account_filter.account, account_filter.owner, None) + .subscribe_with_account_request( + account_filter, + self.event_type_filter.read().await.as_ref(), + ) .unwrap_or_default(); control_sender @@ -386,6 +386,7 @@ impl Clone for YellowstoneGrpc { subscription_handle: self.subscription_handle.clone(), // 共享同一个 Arc> active_subscription: self.active_subscription.clone(), control_tx: self.control_tx.clone(), + event_type_filter: self.event_type_filter.clone(), current_request: self.current_request.clone(), } } diff --git a/src/streaming/yellowstone_sub_system.rs b/src/streaming/yellowstone_sub_system.rs index b572f2d..32c2393 100755 --- a/src/streaming/yellowstone_sub_system.rs +++ b/src/streaming/yellowstone_sub_system.rs @@ -1,6 +1,9 @@ use crate::{ common::AnyResult, - streaming::{grpc::pool::factory, grpc::EventPretty, yellowstone_grpc::YellowstoneGrpc}, + streaming::{ + grpc::{pool::factory, EventPretty}, + yellowstone_grpc::{TransactionFilter, YellowstoneGrpc}, + }, }; use futures::{SinkExt, StreamExt}; use log::error; @@ -39,12 +42,9 @@ impl YellowstoneGrpc { let addrs = vec![SYSTEM_PROGRAM_ID.to_string()]; let account_include = account_include.unwrap_or_default(); let account_exclude = account_exclude.unwrap_or_default(); - let transactions = self.subscription_manager.get_subscribe_request_filter( - account_include, - account_exclude, - addrs, - None, - ); + let tx_filter = + vec![TransactionFilter { account_include, account_exclude, account_required: addrs }]; + let transactions = self.subscription_manager.get_subscribe_request_filter(tx_filter, None); let (mut subscribe_tx, mut stream, _) = self .subscription_manager .subscribe_with_request(transactions, None, None, None)