diff --git a/examples/parse_tx_events.rs b/examples/parse_tx_events.rs index 2384481..91a2007 100644 --- a/examples/parse_tx_events.rs +++ b/examples/parse_tx_events.rs @@ -36,8 +36,9 @@ async fn main() -> Result<()> { /// Get details of a single transaction async fn get_single_transaction_details(signature_str: &str) -> Result<()> { - use solana_sdk::signature::Signature; - use solana_transaction_status::UiTransactionEncoding; + use solana_sdk::{signature::Signature, pubkey::Pubkey, message::compiled_instruction::CompiledInstruction}; + use solana_transaction_status::{UiTransactionEncoding, InnerInstruction, InnerInstructions, UiInstruction}; + use prost_types::Timestamp; let signature = Signature::from_str(signature_str)?; @@ -88,6 +89,96 @@ async fn get_single_transaction_details(signature_str: &str) -> Result<()> { } } } + + // Parse the transaction and extract necessary data + let versioned_tx = match transaction.transaction.transaction.decode() { + Some(tx) => tx, + None => { + println!("Failed to decode transaction"); + return Ok(()); + } + }; + + // Convert inner instructions from meta + let mut inner_instructions_vec: Vec = Vec::new(); + if let Some(meta) = &transaction.transaction.meta { + if let solana_transaction_status::option_serializer::OptionSerializer::Some( + ui_inner_insts, + ) = &meta.inner_instructions + { + for ui_inner in ui_inner_insts { + let mut converted_instructions = Vec::new(); + + for ui_instruction in &ui_inner.instructions { + if let UiInstruction::Compiled(ui_compiled) = ui_instruction { + if let Ok(data) = solana_sdk::bs58::decode(&ui_compiled.data).into_vec() { + let compiled_instruction = CompiledInstruction { + program_id_index: ui_compiled.program_id_index, + accounts: ui_compiled.accounts.to_vec(), + data, + }; + + let inner_instruction = InnerInstruction { + instruction: compiled_instruction, + stack_height: ui_compiled.stack_height, + }; + + converted_instructions.push(inner_instruction); + } + } + } + + let inner_instructions = InnerInstructions { + index: ui_inner.index, + instructions: converted_instructions, + }; + + inner_instructions_vec.push(inner_instructions); + } + } + } + + // Extract address table lookups + let meta = transaction.transaction.meta; + let mut address_table_lookups: Vec = vec![]; + if let Some(meta) = meta { + if let solana_transaction_status::option_serializer::OptionSerializer::Some( + loaded_addresses, + ) = &meta.loaded_addresses + { + address_table_lookups + .reserve(loaded_addresses.writable.len() + loaded_addresses.readonly.len()); + address_table_lookups.extend( + loaded_addresses + .writable + .iter() + .filter_map(|s| s.parse::().ok()) + .chain( + loaded_addresses + .readonly + .iter() + .filter_map(|s| s.parse::().ok()), + ), + ); + } + } + + // Build complete accounts list + let mut accounts = Vec::with_capacity( + versioned_tx.message.static_account_keys().len() + address_table_lookups.len(), + ); + accounts.extend_from_slice(versioned_tx.message.static_account_keys()); + accounts.extend(address_table_lookups); + + let slot = transaction.slot; + let block_time = transaction.block_time.map(|t| Timestamp { seconds: t as i64, nanos: 0 }); + let recv_us = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_micros() as i64; + let bot_wallet = None; + let transaction_index = None; + let protocols = vec![ Protocol::Bonk, Protocol::RaydiumClmm, @@ -96,14 +187,26 @@ async fn get_single_transaction_details(signature_str: &str) -> Result<()> { Protocol::RaydiumCpmm, Protocol::RaydiumAmmV4, ]; - EventParser::parse_encoded_confirmed_transaction_with_status_meta( + + // Create callback + let callback = Arc::new(move |event: DexEvent| { + println!("{:?}\n", event); + }); + + // Call parse_instruction_events_from_versioned_transaction + EventParser::parse_instruction_events_from_versioned_transaction( &protocols, None, + &versioned_tx, signature, - transaction, - Arc::new(move |event: &DexEvent| { - println!("{:?}\n", event); - }), + Some(slot), + block_time, + recv_us, + &accounts, + &inner_instructions_vec, + bot_wallet, + transaction_index, + callback, ) .await?; } diff --git a/src/streaming/common/event_processor.rs b/src/streaming/common/event_processor.rs index 7e7c81e..e415807 100644 --- a/src/streaming/common/event_processor.rs +++ b/src/streaming/common/event_processor.rs @@ -64,7 +64,7 @@ pub async fn process_grpc_transaction( let adapter_callback = create_metrics_callback(callback.clone()); - EventParser::parse_grpc_transaction_owned( + EventParser::parse_grpc_transaction( protocols, event_type_filter, grpc_tx, @@ -123,18 +123,20 @@ pub async fn process_shred_transaction( let recv_us = transaction_with_slot.recv_us; let adapter_callback = create_metrics_callback(callback); + let accounts = tx.message.static_account_keys(); - EventParser::parse_versioned_transaction_owned( + EventParser::parse_instruction_events_from_versioned_transaction( protocols, event_type_filter, - tx, + &tx, signature, Some(slot), None, recv_us, + accounts, + &[], bot_wallet, None, - &[], adapter_callback, ) .await?; diff --git a/src/streaming/event_parser/core/account_event_parser.rs b/src/streaming/event_parser/core/account_event_parser.rs index 0129b4a..5711d04 100644 --- a/src/streaming/event_parser/core/account_event_parser.rs +++ b/src/streaming/event_parser/core/account_event_parser.rs @@ -71,12 +71,12 @@ impl AccountEventParser { if let Some(protocol) = EventDispatcher::match_protocol_by_program_id(&account.owner) { // 检查是否在请求的协议列表中 if protocols.contains(&protocol) { - // 构建临时元数据(具体的event_type会在parser中设置) + // 构建临时元数据(protocol会被dispatcher设置,event_type会在parser中设置) let metadata = EventMetadata { slot: account.slot, signature: account.signature, - protocol: ProtocolType::Common, // 会被具体parser覆盖 - event_type: EventType::default(), // 会被具体parser覆盖 + protocol: ProtocolType::Common, // 会被 EventDispatcher::dispatch_account 设置 + event_type: EventType::default(), // 会被具体 parser 设置 program_id: account.owner, recv_us: account.recv_us, handle_us: elapsed_micros_since(account.recv_us), @@ -144,8 +144,10 @@ impl AccountEventParser { pub fn parse_token_account_event( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::TokenAccount; + let pubkey = account.pubkey; let executable = account.executable; let lamports = account.lamports; @@ -212,8 +214,10 @@ impl AccountEventParser { pub fn parse_nonce_account_event( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::NonceAccount; + if let Ok(info) = parse_nonce(&account.data) { match info { solana_account_decoder::parse_nonce::UiNonceState::Initialized(details) => { diff --git a/src/streaming/event_parser/core/event_parser.rs b/src/streaming/event_parser/core/event_parser.rs index e033a2b..0da6ee1 100644 --- a/src/streaming/event_parser/core/event_parser.rs +++ b/src/streaming/event_parser/core/event_parser.rs @@ -1,66 +1,210 @@ -use crate::streaming::{ - common::SimdUtils, - event_parser::{ - common::{ - filter::EventTypeFilter, - high_performance_clock::{elapsed_micros_since, get_high_perf_clock}, - parse_swap_data_from_next_grpc_instructions, parse_swap_data_from_next_instructions, - EventMetadata, - }, - core::{ - dispatcher::EventDispatcher, - global_state::{ - add_bonk_dev_address, add_dev_address, is_bonk_dev_address_in_signature, - is_dev_address_in_signature, - }, - parser_cache::{ - build_account_pubkeys_with_cache, - get_global_program_ids, - }, - }, - DexEvent, Protocol, +use crate::streaming::event_parser::{ + common::{ + filter::EventTypeFilter, high_performance_clock::elapsed_micros_since, + parse_swap_data_from_next_grpc_instructions, parse_swap_data_from_next_instructions, + EventMetadata, }, + core::{ + dispatcher::EventDispatcher, + global_state::{ + add_bonk_dev_address, add_dev_address, is_bonk_dev_address_in_signature, + is_dev_address_in_signature, + }, + merger_event::merge, + }, + DexEvent, Protocol, }; use prost_types::Timestamp; use solana_sdk::{ - bs58, message::compiled_instruction::CompiledInstruction, pubkey::Pubkey, signature::Signature, + message::compiled_instruction::CompiledInstruction, pubkey::Pubkey, signature::Signature, transaction::VersionedTransaction, }; -use solana_transaction_status::{ - EncodedConfirmedTransactionWithStatusMeta, InnerInstruction, InnerInstructions, UiInstruction, -}; +use solana_transaction_status::InnerInstructions; use std::sync::Arc; use yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionInfo; pub struct EventParser {} impl EventParser { - // 辅助函数:确保 accounts 数组有足够的容量 - fn ensure_accounts_capacity(accounts: &mut Vec, instruction_accounts: &[u8]) { - if let Some(&max_idx) = instruction_accounts.iter().max() { - let required_len = max_idx as usize + 1; - if required_len > accounts.len() { - accounts.resize(required_len, Pubkey::default()); + // ================================================================================================ + // Public API - Entry Points + // ================================================================================================ + + /// Parse transaction from gRPC stream + /// + /// This is the main entry point for parsing transactions received from gRPC streams. + /// It extracts account keys, inner instructions, and delegates to instruction parsing. + pub async fn parse_grpc_transaction( + protocols: &[Protocol], + event_type_filter: Option<&EventTypeFilter>, + grpc_tx: SubscribeUpdateTransactionInfo, + signature: Signature, + slot: Option, + block_time: Option, + recv_us: i64, + bot_wallet: Option, + transaction_index: Option, + callback: Arc, + ) -> anyhow::Result<()> { + // 创建适配器回调,将所有权回调转换为引用回调 + let adapter_callback = Arc::new(move |event: &DexEvent| { + callback(event.clone()); + }); + if let Some(transition) = grpc_tx.transaction { + if let Some(message) = &transition.message { + let mut address_table_lookups: Vec> = vec![]; + let mut inner_instructions: Vec< + yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstructions, + > = vec![]; + + if let Some(meta) = grpc_tx.meta { + inner_instructions = meta.inner_instructions; + address_table_lookups.reserve( + meta.loaded_writable_addresses.len() + meta.loaded_readonly_addresses.len(), + ); + let loaded_writable_addresses = meta.loaded_writable_addresses; + let loaded_readonly_addresses = meta.loaded_readonly_addresses; + address_table_lookups.extend( + loaded_writable_addresses.into_iter().chain(loaded_readonly_addresses), + ); + } + + let mut accounts_bytes: Vec> = + Vec::with_capacity(message.account_keys.len() + address_table_lookups.len()); + accounts_bytes.extend_from_slice(&message.account_keys); + accounts_bytes.extend(address_table_lookups); + // 转换为 Pubkey + let accounts: Vec = accounts_bytes + .iter() + .filter_map(|account| { + if account.len() == 32 { + Some(Pubkey::try_from(account.as_slice()).unwrap_or_default()) + } else { + None + } + }) + .collect(); + // 解析指令事件 + let instructions = &message.instructions; + Self::parse_instruction_events_from_grpc_transaction( + protocols, + event_type_filter, + &instructions, + signature, + slot, + block_time, + recv_us, + &accounts, + &inner_instructions, + bot_wallet, + transaction_index, + adapter_callback, + ) + .await?; } } + + Ok(()) } - // 辅助函数:转换 timestamp 为毫秒 - fn timestamp_to_millis(block_time: Option) -> (Timestamp, i64) { - let timestamp = block_time.unwrap_or(Timestamp { seconds: 0, nanos: 0 }); - let millis = timestamp.seconds * 1000 + (timestamp.nanos as i64) / 1_000_000; - (timestamp, millis) - } - - // 辅助函数:创建回调适配器 - fn create_callback_adapter( + /// Parse transaction from VersionedTransaction + /// + /// This is the entry point for parsing VersionedTransaction objects. + /// It's used when working with RPC responses or historical data. + #[allow(clippy::too_many_arguments)] + pub async fn parse_instruction_events_from_versioned_transaction( + protocols: &[Protocol], + event_type_filter: Option<&EventTypeFilter>, + transaction: &VersionedTransaction, + signature: Signature, + slot: Option, + block_time: Option, + recv_us: i64, + accounts: &[Pubkey], + inner_instructions: &[InnerInstructions], + bot_wallet: Option, + transaction_index: Option, callback: Arc, - ) -> Arc Fn(&'a DexEvent) + Send + Sync> { - Arc::new(move |event: &DexEvent| { + ) -> anyhow::Result<()> { + // 创建适配器回调,将所有权回调转换为引用回调 + let adapter_callback = Arc::new(move |event: &DexEvent| { callback(event.clone()); - }) + }); + // 获取交易的指令和账户 + let compiled_instructions = transaction.message.instructions(); + let mut accounts: Vec = accounts.to_vec(); + // 检查交易中是否包含程序 + let has_program = accounts + .iter() + .any(|account| Self::should_handle(protocols, event_type_filter, account)); + if has_program { + // 解析每个指令 + for (index, instruction) in compiled_instructions.iter().enumerate() { + if let Some(program_id) = accounts.get(instruction.program_id_index as usize) { + let program_id = *program_id; // 克隆程序ID,避免借用冲突 + let inner_instructions = inner_instructions + .iter() + .find(|inner_instruction| inner_instruction.index == index as u8); + if Self::should_handle(protocols, event_type_filter, &program_id) { + let max_idx = instruction.accounts.iter().max().unwrap_or(&0); + // 补齐accounts(使用Pubkey::default()) + if *max_idx as usize >= accounts.len() { + accounts.resize(*max_idx as usize + 1, Pubkey::default()); + } + Self::parse_events_from_instruction( + protocols, + event_type_filter, + instruction, + &accounts, + signature, + slot.unwrap_or(0), + block_time, + recv_us, + index as i64, + None, + bot_wallet, + transaction_index, + inner_instructions, + adapter_callback.clone(), + )?; + } + // Immediately process inner instructions for correct ordering + if let Some(inner_instructions) = inner_instructions { + for (inner_index, inner_instruction) in + inner_instructions.instructions.iter().enumerate() + { + Self::parse_events_from_instruction( + protocols, + event_type_filter, + &inner_instruction.instruction, + &accounts, + signature, + slot.unwrap_or(0), + block_time, + recv_us, + index as i64, + Some(inner_index as i64), + bot_wallet, + transaction_index, + Some(&inner_instructions), + adapter_callback.clone(), + )?; + } + } + } + } + } + Ok(()) } + // ================================================================================================ + // gRPC Transaction Processing + // ================================================================================================ + + /// Parse instruction events from gRPC transaction format + /// + /// Iterates through all instructions in a gRPC transaction, checks if they should be handled, + /// and delegates to instruction-level parsing for both outer and inner instructions. #[allow(clippy::too_many_arguments)] async fn parse_instruction_events_from_grpc_transaction( protocols: &[Protocol], @@ -90,8 +234,11 @@ impl EventParser { let inner_instructions = inner_instructions .iter() .find(|inner_instruction| inner_instruction.index == index as u32); + let max_idx = instruction.accounts.iter().max().unwrap_or(&0); // 补齐accounts(使用Pubkey::default()) - Self::ensure_accounts_capacity(&mut accounts, &instruction.accounts); + if *max_idx as usize >= accounts.len() { + accounts.resize(*max_idx as usize + 1, Pubkey::default()); + } if Self::should_handle(protocols, event_type_filter, &program_id) { Self::parse_events_from_grpc_instruction( protocols, @@ -147,431 +294,10 @@ impl EventParser { Ok(()) } - /// 从VersionedTransaction中解析指令事件的通用方法 - #[allow(clippy::too_many_arguments)] - async fn parse_instruction_events_from_versioned_transaction( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - transaction: &VersionedTransaction, - signature: Signature, - slot: Option, - block_time: Option, - recv_us: i64, - accounts: &[Pubkey], - inner_instructions: &[InnerInstructions], - bot_wallet: Option, - transaction_index: Option, - callback: Arc Fn(&'a DexEvent) + Send + Sync>, - ) -> anyhow::Result<()> { - // 获取交易的指令和账户 - let compiled_instructions = transaction.message.instructions(); - let mut accounts: Vec = accounts.to_vec(); - // 检查交易中是否包含程序 - let has_program = accounts - .iter() - .any(|account| Self::should_handle(protocols, event_type_filter, account)); - if has_program { - // 解析每个指令 - for (index, instruction) in compiled_instructions.iter().enumerate() { - if let Some(program_id) = accounts.get(instruction.program_id_index as usize) { - let program_id = *program_id; // 克隆程序ID,避免借用冲突 - let inner_instructions = inner_instructions - .iter() - .find(|inner_instruction| inner_instruction.index == index as u8); - // 补齐accounts(使用Pubkey::default()) - Self::ensure_accounts_capacity(&mut accounts, &instruction.accounts); - if Self::should_handle(protocols, event_type_filter, &program_id) { - Self::parse_events_from_instruction( - protocols, - event_type_filter, - instruction, - &accounts, - signature, - slot.unwrap_or(0), - block_time, - recv_us, - index as i64, - None, - bot_wallet, - transaction_index, - inner_instructions, - callback.clone(), - )?; - } - // Immediately process inner instructions for correct ordering - if let Some(inner_instructions) = inner_instructions { - for (inner_index, inner_instruction) in - inner_instructions.instructions.iter().enumerate() - { - Self::parse_events_from_instruction( - protocols, - event_type_filter, - &inner_instruction.instruction, - &accounts, - signature, - slot.unwrap_or(0), - block_time, - recv_us, - index as i64, - Some(inner_index as i64), - bot_wallet, - transaction_index, - Some(&inner_instructions), - callback.clone(), - )?; - } - } - } - } - } - Ok(()) - } - - pub async fn parse_versioned_transaction_owned( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - versioned_tx: VersionedTransaction, - signature: Signature, - slot: Option, - block_time: Option, - recv_us: i64, - bot_wallet: Option, - transaction_index: Option, - inner_instructions: &[InnerInstructions], - callback: Arc, - ) -> anyhow::Result<()> { - let adapter_callback = Self::create_callback_adapter(callback); - let accounts = versioned_tx.message.static_account_keys(); - Self::parse_instruction_events_from_versioned_transaction( - protocols, - event_type_filter, - &versioned_tx, - signature, - slot, - block_time, - recv_us, - accounts, - inner_instructions, - bot_wallet, - transaction_index, - adapter_callback, - ) - .await - } - - pub async fn parse_grpc_transaction_owned( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - grpc_tx: SubscribeUpdateTransactionInfo, - signature: Signature, - slot: Option, - block_time: Option, - recv_us: i64, - bot_wallet: Option, - transaction_index: Option, - callback: Arc, - ) -> anyhow::Result<()> { - let adapter_callback = Self::create_callback_adapter(callback); - Self::parse_grpc_transaction( - protocols, - event_type_filter, - grpc_tx, - signature, - slot, - block_time, - recv_us, - bot_wallet, - transaction_index, - adapter_callback, - ) - .await - } - - async fn parse_grpc_transaction( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - grpc_tx: SubscribeUpdateTransactionInfo, - signature: Signature, - slot: Option, - block_time: Option, - recv_us: i64, - bot_wallet: Option, - transaction_index: Option, - callback: Arc Fn(&'a DexEvent) + Send + Sync>, - ) -> anyhow::Result<()> { - if let Some(transition) = grpc_tx.transaction { - if let Some(message) = &transition.message { - let mut address_table_lookups: Vec> = vec![]; - let mut inner_instructions: Vec< - yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstructions, - > = vec![]; - - if let Some(meta) = grpc_tx.meta { - inner_instructions = meta.inner_instructions; - address_table_lookups.reserve( - meta.loaded_writable_addresses.len() + meta.loaded_readonly_addresses.len(), - ); - let loaded_writable_addresses = meta.loaded_writable_addresses; - let loaded_readonly_addresses = meta.loaded_readonly_addresses; - address_table_lookups.extend( - loaded_writable_addresses.into_iter().chain(loaded_readonly_addresses), - ); - } - - let mut accounts_bytes: Vec> = - Vec::with_capacity(message.account_keys.len() + address_table_lookups.len()); - accounts_bytes.extend_from_slice(&message.account_keys); - accounts_bytes.extend(address_table_lookups); - // 转换为 Pubkey - let accounts: Vec = accounts_bytes - .iter() - .filter_map(|account| { - if account.len() == 32 { - Some(Pubkey::try_from(account.as_slice()).unwrap_or_default()) - } else { - None - } - }) - .collect(); - // 解析指令事件 - let instructions = &message.instructions; - Self::parse_instruction_events_from_grpc_transaction( - protocols, - event_type_filter, - &instructions, - signature, - slot, - block_time, - recv_us, - &accounts, - &inner_instructions, - bot_wallet, - transaction_index, - callback, - ) - .await?; - } - } - - Ok(()) - } - - pub async fn parse_encoded_confirmed_transaction_with_status_meta( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - signature: Signature, - transaction: EncodedConfirmedTransactionWithStatusMeta, - callback: Arc Fn(&'a DexEvent) + Send + Sync>, - ) -> anyhow::Result<()> { - let versioned_tx = match transaction.transaction.transaction.decode() { - Some(tx) => tx, - None => { - return Ok(()); - } - }; - let mut inner_instructions_vec: Vec = Vec::new(); - if let Some(meta) = &transaction.transaction.meta { - // 从meta中获取inner_instructions,处理OptionSerializer类型 - if let solana_transaction_status::option_serializer::OptionSerializer::Some( - ui_inner_insts, - ) = &meta.inner_instructions - { - // 将UiInnerInstructions转换为InnerInstructions - for ui_inner in ui_inner_insts { - let mut converted_instructions = Vec::new(); - - // 转换每个UiInstruction为InnerInstruction - for ui_instruction in &ui_inner.instructions { - if let UiInstruction::Compiled(ui_compiled) = ui_instruction { - // 解码base58编码的data - if let Ok(data) = bs58::decode(&ui_compiled.data).into_vec() { - // base64解码 - let compiled_instruction = CompiledInstruction { - program_id_index: ui_compiled.program_id_index, - accounts: ui_compiled.accounts.to_vec(), - data, - }; - - let inner_instruction = InnerInstruction { - instruction: compiled_instruction, - stack_height: ui_compiled.stack_height, - }; - - converted_instructions.push(inner_instruction); - } - } - } - - let inner_instructions = InnerInstructions { - index: ui_inner.index, - instructions: converted_instructions, - }; - - inner_instructions_vec.push(inner_instructions); - } - } - } - let inner_instructions: &[InnerInstructions] = &inner_instructions_vec; - - let meta = transaction.transaction.meta; - let mut address_table_lookups: Vec = vec![]; - if let Some(meta) = meta { - if let solana_transaction_status::option_serializer::OptionSerializer::Some( - loaded_addresses, - ) = &meta.loaded_addresses - { - address_table_lookups - .reserve(loaded_addresses.writable.len() + loaded_addresses.readonly.len()); - address_table_lookups.extend( - loaded_addresses - .writable - .iter() - .filter_map(|s| s.parse::().ok()) - .chain( - loaded_addresses - .readonly - .iter() - .filter_map(|s| s.parse::().ok()), - ), - ); - } - } - let mut accounts = Vec::with_capacity( - versioned_tx.message.static_account_keys().len() + address_table_lookups.len(), - ); - accounts.extend_from_slice(versioned_tx.message.static_account_keys()); - accounts.extend(address_table_lookups); - - let slot = transaction.slot; - let block_time = transaction.block_time.map(|t| Timestamp { seconds: t as i64, nanos: 0 }); - let recv_us = get_high_perf_clock(); - let bot_wallet = None; - let transaction_index = None; - // 解析指令事件 - Self::parse_instruction_events_from_versioned_transaction( - protocols, - event_type_filter, - &versioned_tx, - signature, - Some(slot), - block_time, - recv_us, - &accounts, - inner_instructions, - bot_wallet, - transaction_index, - callback, - ) - .await - } - - /// 从指令中解析事件 - #[allow(clippy::too_many_arguments)] - fn parse_events_from_instruction( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - instruction: &CompiledInstruction, - accounts: &[Pubkey], - signature: Signature, - slot: u64, - block_time: Option, - recv_us: i64, - outer_index: i64, - inner_index: Option, - bot_wallet: Option, - transaction_index: Option, - inner_instructions: Option<&InnerInstructions>, - callback: Arc Fn(&'a DexEvent) + Send + Sync>, - ) -> anyhow::Result<()> { - // 验证 program_id - let program_id = match Self::validate_program_id( - instruction.program_id_index as u32, - accounts, - protocols, - event_type_filter, - ) { - Some(id) => id, - None => return Ok(()), - }; - - // 验证并构建账户公钥列表 - let account_pubkeys = match Self::validate_and_build_account_pubkeys(&instruction.accounts, accounts) { - Some(pubkeys) => pubkeys, - None => return Ok(()), - }; - - // 提取判别器 - let (instruction_discriminator, instruction_data) = - match Self::extract_instruction_discriminator(&instruction.data) { - Some(result) => result, - None => return Ok(()), - }; - - // 准备 inner instruction 数据(如果有) - let (inner_discriminator, inner_data) = if let Some(inner_instructions_ref) = inner_instructions { - // 尝试从第一个 inner instruction 获取数据 - if let Some(first_inner) = inner_instructions_ref.instructions.first() { - let inner_inst_data = &first_inner.instruction.data; - if inner_inst_data.len() >= 24 { // 16 bytes discriminator + data - (Some(&inner_inst_data[0..16]), Some(&inner_inst_data[16..])) - } else { - (None, None) - } - } else { - (None, None) - } - } else { - (None, None) - }; - - // 创建元数据(使用默认值,dispatcher会返回带正确元数据的事件) - let (timestamp, block_time_ms) = Self::timestamp_to_millis(block_time); - let metadata = EventMetadata::new( - signature, - slot, - timestamp.seconds, - block_time_ms, - crate::streaming::event_parser::common::ProtocolType::Common, - crate::streaming::event_parser::common::EventType::default(), - program_id, - outer_index, - inner_index, - recv_us, - transaction_index, - ); - - // 使用 dispatcher 解析事件 - if let Some(mut event) = EventDispatcher::dispatch_multi( - protocols, - &program_id, - instruction_discriminator, - inner_discriminator, - instruction_data, - inner_data, - &account_pubkeys, - metadata, - ) { - // 解析 swap_data(如果需要) - if let Some(inner_instructions_ref) = inner_instructions { - if event.metadata().swap_data.is_none() { - if let Some(swap_data) = parse_swap_data_from_next_instructions( - &event, - inner_instructions_ref, - inner_index.unwrap_or(-1_i64) as i8, - accounts, - ) { - event.metadata_mut().set_swap_data(swap_data); - } - } - } - - // 完成事件处理 - Self::finalize_event(event, recv_us, bot_wallet, &callback); - } - Ok(()) - } - - /// 从指令中解析事件 (gRPC版本) + /// Parse events from gRPC instruction + /// + /// Core parsing logic for a single gRPC instruction. Extracts discriminator, dispatches + /// to protocol-specific parsers, handles inner instructions, and processes swap data. #[allow(clippy::too_many_arguments)] fn parse_events_from_grpc_instruction( protocols: &[Protocol], @@ -589,56 +315,48 @@ impl EventParser { inner_instructions: Option<&yellowstone_grpc_proto::prelude::InnerInstructions>, callback: Arc Fn(&'a DexEvent) + Send + Sync>, ) -> anyhow::Result<()> { - // 验证 program_id - let program_id = match Self::validate_program_id( - instruction.program_id_index, - accounts, - protocols, - event_type_filter, - ) { - Some(id) => id, + // 添加边界检查以防止越界访问 + let program_id_index = instruction.program_id_index as usize; + if program_id_index >= accounts.len() { + return Ok(()); + } + let program_id = accounts[program_id_index]; + if !Self::should_handle(protocols, event_type_filter, &program_id) { + return Ok(()); + } + + // 使用 EventDispatcher 匹配协议 + let protocol = match EventDispatcher::match_protocol_by_program_id(&program_id) { + Some(p) => p, None => return Ok(()), }; - // 验证并构建账户公钥列表 - let account_pubkeys = match Self::validate_and_build_account_pubkeys(&instruction.accounts, accounts) { - Some(pubkeys) => pubkeys, - None => return Ok(()), - }; + // 检查指令数据长度(至少需要 8 字节的 discriminator) + if instruction.data.len() < 8 { + return Ok(()); + } - // 提取判别器 - let (instruction_discriminator, instruction_data) = - match Self::extract_instruction_discriminator(&instruction.data) { - Some(result) => result, - None => return Ok(()), - }; + // 提取 discriminator 和数据 + let instruction_discriminator = &instruction.data[..8]; + let instruction_data = &instruction.data[8..]; - // 准备 inner instruction 数据(如果有) - let (inner_discriminator, inner_data) = if let Some(inner_instructions_ref) = inner_instructions { - // 尝试从第一个 inner instruction 获取数据 - if let Some(first_inner) = inner_instructions_ref.instructions.first() { - let inner_inst_data = &first_inner.data; - if inner_inst_data.len() >= 24 { // 16 bytes discriminator + data - (Some(&inner_inst_data[0..16]), Some(&inner_inst_data[16..])) - } else { - (None, None) - } - } else { - (None, None) - } - } else { - (None, None) - }; + // 构建账户公钥列表 + let account_pubkeys: Vec = instruction + .accounts + .iter() + .filter_map(|&idx| accounts.get(idx as usize).copied()) + .collect(); - // 创建元数据(使用默认值,dispatcher会返回带正确元数据的事件) - let (timestamp, block_time_ms) = Self::timestamp_to_millis(block_time); + // 创建元数据 + let timestamp = block_time.unwrap_or(Timestamp { seconds: 0, nanos: 0 }); + let block_time_ms = timestamp.seconds * 1000 + (timestamp.nanos as i64) / 1_000_000; let metadata = EventMetadata::new( signature, slot, timestamp.seconds, block_time_ms, - crate::streaming::event_parser::common::ProtocolType::Common, - crate::streaming::event_parser::common::EventType::default(), + Default::default(), // protocol will be set by dispatcher + Default::default(), // event_type will be set by dispatcher program_id, outer_index, inner_index, @@ -646,105 +364,279 @@ impl EventParser { transaction_index, ); - // 使用 dispatcher 解析事件 - if let Some(mut event) = EventDispatcher::dispatch_multi( - protocols, - &program_id, + // 使用 EventDispatcher 解析 instruction 事件 + let mut event = match EventDispatcher::dispatch_instruction( + protocol.clone(), instruction_discriminator, - inner_discriminator, instruction_data, - inner_data, &account_pubkeys, - metadata, + metadata.clone(), ) { - // 解析 swap_data(如果需要) - if let Some(inner_instructions_ref) = inner_instructions { - if event.metadata().swap_data.is_none() { - if let Some(swap_data) = parse_swap_data_from_next_grpc_instructions( - &event, - inner_instructions_ref, - inner_index.unwrap_or(-1_i64) as i8, - accounts, - ) { - event.metadata_mut().set_swap_data(swap_data); + Some(e) => e, + None => return Ok(()), + }; + + // 处理 inner instructions + let mut inner_instruction_event: Option = None; + if let Some(inner_instructions_ref) = inner_instructions { + // 并行执行两个任务: 解析 inner event 和提取 swap_data + let (inner_event_result, swap_data_result) = std::thread::scope(|s| { + let inner_event_handle = s.spawn(|| { + for inner_instruction in inner_instructions_ref.instructions.iter() { + let inner_data = &inner_instruction.data; + // 检查长度(需要 16 字节的 discriminator) + if inner_data.len() < 16 { + continue; + } + let inner_discriminator = &inner_data[..16]; + let inner_instruction_data = &inner_data[16..]; + + if let Some(inner_event) = EventDispatcher::dispatch_inner_instruction( + protocol.clone(), + inner_discriminator, + inner_instruction_data, + metadata.clone(), + ) { + return Some(inner_event); + } } - } + None + }); + + let swap_data_handle = s.spawn(|| { + if event.metadata().swap_data.is_none() { + parse_swap_data_from_next_grpc_instructions( + &event, + inner_instructions_ref, + inner_index.unwrap_or(-1_i64) as i8, + accounts, + ) + } else { + None + } + }); + + // 等待两个任务完成 + (inner_event_handle.join().unwrap(), swap_data_handle.join().unwrap()) + }); + + inner_instruction_event = inner_event_result; + if let Some(swap_data) = swap_data_result { + event.metadata_mut().set_swap_data(swap_data); } - - // 完成事件处理 - Self::finalize_event(event, recv_us, bot_wallet, &callback); } - Ok(()) - } - fn should_handle( - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - program_id: &Pubkey, - ) -> bool { - get_global_program_ids(protocols, event_type_filter).contains(program_id) - } - - // 辅助函数:设置 swap_data 的 from_amount 和 to_amount - fn set_swap_amounts( - swap_data: &mut crate::streaming::event_parser::common::SwapData, - from_amount: u64, - to_amount: u64, - ) { - swap_data.from_amount = from_amount; - swap_data.to_amount = to_amount; - } - - // 辅助函数:验证 program_id 是否在 accounts 范围内并且应该被处理 - fn validate_program_id( - program_id_index: u32, - accounts: &[Pubkey], - protocols: &[Protocol], - event_type_filter: Option<&EventTypeFilter>, - ) -> Option { - let index = program_id_index as usize; - if index >= accounts.len() { - return None; + // 特殊处理: PumpFun MIGRATE 指令需要 inner instruction data + if matches!(protocol, Protocol::PumpFun) { + const PUMPFUN_MIGRATE_IX: &[u8] = &[155, 234, 231, 146, 236, 158, 162, 30]; + if instruction_discriminator == PUMPFUN_MIGRATE_IX && inner_instruction_event.is_none() + { + return Ok(()); + } } - let program_id = accounts[index]; - if Self::should_handle(protocols, event_type_filter, &program_id) { - Some(program_id) - } else { - None - } - } - // 辅助函数:验证并构建账户公钥列表 - fn validate_and_build_account_pubkeys( - instruction_accounts: &[u8], - accounts: &[Pubkey], - ) -> Option> { - if !SimdUtils::validate_account_indices_simd(instruction_accounts, accounts.len()) { - return None; + // 合并事件 + if let Some(inner_instruction_event) = inner_instruction_event { + merge(&mut event, inner_instruction_event); } - Some(build_account_pubkeys_with_cache(instruction_accounts, accounts)) - } - // 辅助函数:提取指令数据的判别器 - fn extract_instruction_discriminator(instruction_data: &[u8]) -> Option<(&[u8], &[u8])> { - if instruction_data.len() < 8 { - return None; - } - Some((&instruction_data[0..8], &instruction_data[8..])) - } - - // 辅助函数:处理事件的最后步骤(设置时间、处理事件、调用回调) - fn finalize_event( - mut event: DexEvent, - recv_us: i64, - bot_wallet: Option, - callback: &Arc Fn(&'a DexEvent) + Send + Sync>, - ) { + // 设置处理时间(使用高性能时钟) event.metadata_mut().handle_us = elapsed_micros_since(recv_us); event = Self::process_event(event, bot_wallet); callback(&event); + + Ok(()) } + // ================================================================================================ + // Standard Instruction Processing + // ================================================================================================ + + /// Parse events from standard Solana instruction + /// + /// Similar to gRPC instruction parsing but works with standard CompiledInstruction format. + /// Used when parsing VersionedTransaction or RPC data. + #[allow(clippy::too_many_arguments)] + fn parse_events_from_instruction( + protocols: &[Protocol], + event_type_filter: Option<&EventTypeFilter>, + instruction: &CompiledInstruction, + accounts: &[Pubkey], + signature: Signature, + slot: u64, + block_time: Option, + recv_us: i64, + outer_index: i64, + inner_index: Option, + bot_wallet: Option, + transaction_index: Option, + inner_instructions: Option<&InnerInstructions>, + callback: Arc Fn(&'a DexEvent) + Send + Sync>, + ) -> anyhow::Result<()> { + // 添加边界检查以防止越界访问 + let program_id_index = instruction.program_id_index as usize; + if program_id_index >= accounts.len() { + return Ok(()); + } + let program_id = accounts[program_id_index]; + if !Self::should_handle(protocols, event_type_filter, &program_id) { + return Ok(()); + } + + // 使用 EventDispatcher 匹配协议 + let protocol = match EventDispatcher::match_protocol_by_program_id(&program_id) { + Some(p) => p, + None => return Ok(()), + }; + + // 检查指令数据长度(至少需要 8 字节的 discriminator) + if instruction.data.len() < 8 { + return Ok(()); + } + + // 提取 discriminator 和数据 + let instruction_discriminator = &instruction.data[..8]; + let instruction_data = &instruction.data[8..]; + + // 构建账户公钥列表 + let account_pubkeys: Vec = instruction + .accounts + .iter() + .filter_map(|&idx| accounts.get(idx as usize).copied()) + .collect(); + + // 创建元数据 + let timestamp = block_time.unwrap_or(Timestamp { seconds: 0, nanos: 0 }); + let block_time_ms = timestamp.seconds * 1000 + (timestamp.nanos as i64) / 1_000_000; + let metadata = EventMetadata::new( + signature, + slot, + timestamp.seconds, + block_time_ms, + Default::default(), // protocol will be set by dispatcher + Default::default(), // event_type will be set by dispatcher + program_id, + outer_index, + inner_index, + recv_us, + transaction_index, + ); + + // 使用 EventDispatcher 解析 instruction 事件 + let mut event = match EventDispatcher::dispatch_instruction( + protocol.clone(), + instruction_discriminator, + instruction_data, + &account_pubkeys, + metadata.clone(), + ) { + Some(e) => e, + None => return Ok(()), + }; + + // 处理 inner instructions + let mut inner_instruction_event: Option = None; + if let Some(inner_instructions_ref) = inner_instructions { + // 并行执行两个任务: 解析 inner event 和提取 swap_data + let (inner_event_result, swap_data_result) = std::thread::scope(|s| { + let inner_event_handle = s.spawn(|| { + for inner_instruction in inner_instructions_ref.instructions.iter() { + let inner_data = &inner_instruction.instruction.data; + // 检查长度(需要 16 字节的 discriminator) + if inner_data.len() < 16 { + continue; + } + let inner_discriminator = &inner_data[..16]; + let inner_instruction_data = &inner_data[16..]; + + if let Some(inner_event) = EventDispatcher::dispatch_inner_instruction( + protocol.clone(), + inner_discriminator, + inner_instruction_data, + metadata.clone(), + ) { + return Some(inner_event); + } + } + None + }); + + let swap_data_handle = s.spawn(|| { + if event.metadata().swap_data.is_none() { + parse_swap_data_from_next_instructions( + &event, + inner_instructions_ref, + inner_index.unwrap_or(-1_i64) as i8, + accounts, + ) + } else { + None + } + }); + + // 等待两个任务完成 + (inner_event_handle.join().unwrap(), swap_data_handle.join().unwrap()) + }); + + inner_instruction_event = inner_event_result; + if let Some(swap_data) = swap_data_result { + event.metadata_mut().set_swap_data(swap_data); + } + } + + // 特殊处理: PumpFun MIGRATE 指令需要 inner instruction data + if matches!(protocol, Protocol::PumpFun) { + const PUMPFUN_MIGRATE_IX: &[u8] = &[155, 234, 231, 146, 236, 158, 162, 30]; + if instruction_discriminator == PUMPFUN_MIGRATE_IX && inner_instruction_event.is_none() + { + return Ok(()); + } + } + + // 合并事件 + if let Some(inner_instruction_event) = inner_instruction_event { + merge(&mut event, inner_instruction_event); + } + + // 设置处理时间(使用高性能时钟) + event.metadata_mut().handle_us = elapsed_micros_since(recv_us); + event = Self::process_event(event, bot_wallet); + callback(&event); + + Ok(()) + } + + // ================================================================================================ + // Helper Functions + // ================================================================================================ + + /// Check if instruction should be processed based on protocol filter + /// + /// Determines whether a program_id matches any of the protocols we're interested in. + fn should_handle( + protocols: &[Protocol], + _event_type_filter: Option<&EventTypeFilter>, + program_id: &Pubkey, + ) -> bool { + // 使用 EventDispatcher 来匹配协议 + if let Some(protocol) = EventDispatcher::match_protocol_by_program_id(program_id) { + protocols.contains(&protocol) + } else { + false + } + } + + // ================================================================================================ + // Event Post-Processing + // ================================================================================================ + + /// Process and enrich parsed event with additional context + /// + /// Handles protocol-specific post-processing: + /// - PumpFun: Tracks dev addresses and marks dev trades + /// - PumpSwap: Fills swap data amounts + /// - Bonk: Tracks pool creators and marks dev trades + /// - General: Marks bot wallet trades fn process_event(event: DexEvent, bot_wallet: Option) -> DexEvent { let signature = event.metadata().signature; // Copy the signature to avoid borrowing issues match event { @@ -763,32 +655,30 @@ impl EventParser { trade_info.is_bot = Some(trade_info.user) == bot_wallet; if let Some(swap_data) = trade_info.metadata.swap_data.as_mut() { - let (from, to) = if trade_info.is_buy { - (trade_info.sol_amount, trade_info.token_amount) + swap_data.from_amount = if trade_info.is_buy { + trade_info.sol_amount } else { - (trade_info.token_amount, trade_info.sol_amount) + trade_info.token_amount + }; + swap_data.to_amount = if trade_info.is_buy { + trade_info.token_amount + } else { + trade_info.sol_amount }; - Self::set_swap_amounts(swap_data, from, to); } DexEvent::PumpFunTradeEvent(trade_info) } DexEvent::PumpSwapBuyEvent(mut trade_info) => { if let Some(swap_data) = trade_info.metadata.swap_data.as_mut() { - Self::set_swap_amounts( - swap_data, - trade_info.user_quote_amount_in, - trade_info.base_amount_out, - ); + swap_data.from_amount = trade_info.user_quote_amount_in; + swap_data.to_amount = trade_info.base_amount_out; } DexEvent::PumpSwapBuyEvent(trade_info) } DexEvent::PumpSwapSellEvent(mut trade_info) => { if let Some(swap_data) = trade_info.metadata.swap_data.as_mut() { - Self::set_swap_amounts( - swap_data, - trade_info.base_amount_in, - trade_info.user_quote_amount_out, - ); + swap_data.from_amount = trade_info.base_amount_in; + swap_data.to_amount = trade_info.user_quote_amount_out; } DexEvent::PumpSwapSellEvent(trade_info) } diff --git a/src/streaming/event_parser/protocols/bonk/types.rs b/src/streaming/event_parser/protocols/bonk/types.rs index 5386852..a2ddbef 100755 --- a/src/streaming/event_parser/protocols/bonk/types.rs +++ b/src/streaming/event_parser/protocols/bonk/types.rs @@ -4,7 +4,7 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, + common::{EventMetadata, EventType}, protocols::bonk::{ BonkGlobalConfigAccountEvent, BonkPlatformConfigAccountEvent, BonkPoolStateAccountEvent, }, @@ -132,7 +132,9 @@ pub fn pool_state_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..POOL_STATE_SIZE]).ok() } -pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountBonkPoolState; + if account.data.len() < POOL_STATE_SIZE + 8 { return None; } @@ -182,8 +184,10 @@ pub fn global_config_decode(data: &[u8]) -> Option { pub fn global_config_parser( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::AccountBonkGlobalConfig; + if account.data.len() < GLOBAL_CONFIG_SIZE + 8 { return None; } @@ -228,8 +232,10 @@ pub fn platform_config_decode(data: &[u8]) -> Option { pub fn platform_config_parser( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::AccountBonkPlatformConfig; + if account.data.len() < PLATFORM_CONFIG_SIZE + 8 { return None; } diff --git a/src/streaming/event_parser/protocols/pumpfun/types.rs b/src/streaming/event_parser/protocols/pumpfun/types.rs index 15a0e23..0805fb4 100644 --- a/src/streaming/event_parser/protocols/pumpfun/types.rs +++ b/src/streaming/event_parser/protocols/pumpfun/types.rs @@ -4,7 +4,7 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, + common::{EventMetadata, EventType}, protocols::pumpfun::{PumpFunBondingCurveAccountEvent, PumpFunGlobalAccountEvent}, DexEvent, }, @@ -33,8 +33,10 @@ pub fn bonding_curve_decode(data: &[u8]) -> Option { pub fn bonding_curve_parser( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::AccountPumpFunBondingCurve; + if account.data.len() < BONDING_CURVE_SIZE + 8 { return None; } @@ -81,7 +83,9 @@ pub fn global_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..GLOBAL_SIZE]).ok() } -pub fn global_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn global_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountPumpFunGlobal; + if account.data.len() < GLOBAL_SIZE + 8 { return None; } diff --git a/src/streaming/event_parser/protocols/pumpswap/types.rs b/src/streaming/event_parser/protocols/pumpswap/types.rs index e66d007..afa99bc 100644 --- a/src/streaming/event_parser/protocols/pumpswap/types.rs +++ b/src/streaming/event_parser/protocols/pumpswap/types.rs @@ -4,7 +4,7 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, + common::{EventMetadata, EventType}, protocols::pumpswap::{PumpSwapGlobalConfigAccountEvent, PumpSwapPoolAccountEvent}, DexEvent, }, @@ -33,8 +33,10 @@ pub fn global_config_decode(data: &[u8]) -> Option { pub fn global_config_parser( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::AccountPumpSwapGlobalConfig; + if account.data.len() < GLOBAL_CONFIG_SIZE + 8 { return None; } @@ -76,7 +78,9 @@ pub fn pool_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..POOL_SIZE]).ok() } -pub fn pool_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn pool_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountPumpSwapPool; + if account.data.len() < POOL_SIZE + 8 { return None; } diff --git a/src/streaming/event_parser/protocols/raydium_amm_v4/types.rs b/src/streaming/event_parser/protocols/raydium_amm_v4/types.rs index ba8d7ce..073dac6 100644 --- a/src/streaming/event_parser/protocols/raydium_amm_v4/types.rs +++ b/src/streaming/event_parser/protocols/raydium_amm_v4/types.rs @@ -4,7 +4,8 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, protocols::raydium_amm_v4::RaydiumAmmV4AmmInfoAccountEvent, + common::{EventMetadata, EventType}, + protocols::raydium_amm_v4::RaydiumAmmV4AmmInfoAccountEvent, DexEvent, }, grpc::AccountPretty, @@ -86,7 +87,9 @@ pub fn amm_info_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..AMM_INFO_SIZE]).ok() } -pub fn amm_info_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn amm_info_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountRaydiumAmmV4AmmInfo; + if account.data.len() < AMM_INFO_SIZE { return None; } diff --git a/src/streaming/event_parser/protocols/raydium_clmm/types.rs b/src/streaming/event_parser/protocols/raydium_clmm/types.rs index a3731b5..5f3695c 100644 --- a/src/streaming/event_parser/protocols/raydium_clmm/types.rs +++ b/src/streaming/event_parser/protocols/raydium_clmm/types.rs @@ -4,7 +4,7 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, + common::{EventMetadata, EventType}, protocols::raydium_clmm::{ RaydiumClmmAmmConfigAccountEvent, RaydiumClmmPoolStateAccountEvent, RaydiumClmmTickArrayStateAccountEvent, @@ -37,7 +37,9 @@ pub fn amm_config_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..AMM_CONFIG_SIZE]).ok() } -pub fn amm_config_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn amm_config_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountRaydiumClmmAmmConfig; + if account.data.len() < AMM_CONFIG_SIZE + 8 { return None; } @@ -122,7 +124,9 @@ pub fn pool_state_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..POOL_STATE_SIZE]).ok() } -pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountRaydiumClmmPoolState; + if account.data.len() < POOL_STATE_SIZE + 8 { return None; } @@ -202,8 +206,10 @@ pub fn tick_array_state_decode(data: &[u8]) -> Option { pub fn tick_array_state_parser( account: &AccountPretty, - metadata: EventMetadata, + mut metadata: EventMetadata, ) -> Option { + metadata.event_type = EventType::AccountRaydiumClmmTickArrayState; + if account.data.len() < TICK_ARRAY_STATE_SIZE + 8 { return None; } diff --git a/src/streaming/event_parser/protocols/raydium_cpmm/types.rs b/src/streaming/event_parser/protocols/raydium_cpmm/types.rs index 21601a2..5b8e9fe 100644 --- a/src/streaming/event_parser/protocols/raydium_cpmm/types.rs +++ b/src/streaming/event_parser/protocols/raydium_cpmm/types.rs @@ -4,7 +4,7 @@ use solana_sdk::pubkey::Pubkey; use crate::streaming::{ event_parser::{ - common::EventMetadata, + common::{EventMetadata, EventType}, protocols::raydium_cpmm::{ RaydiumCpmmAmmConfigAccountEvent, RaydiumCpmmPoolStateAccountEvent, }, @@ -36,7 +36,9 @@ pub fn amm_config_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..AMM_CONFIG_SIZE]).ok() } -pub fn amm_config_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn amm_config_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountRaydiumCpmmAmmConfig; + if account.data.len() < AMM_CONFIG_SIZE + 8 { return None; } @@ -91,7 +93,9 @@ pub fn pool_state_decode(data: &[u8]) -> Option { borsh::from_slice::(&data[..POOL_STATE_SIZE]).ok() } -pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option { +pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option { + metadata.event_type = EventType::AccountRaydiumCpmmPoolState; + if account.data.len() < POOL_STATE_SIZE + 8 { return None; }