refactor: unify event parsing API with callback ownership transfer

This commit is contained in:
ysq
2025-11-04 22:49:40 +08:00
parent 4c0e29865f
commit 819b71ce1f
10 changed files with 665 additions and 639 deletions
+6 -4
View File
@@ -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?;
@@ -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<DexEvent> {
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<DexEvent> {
metadata.event_type = EventType::NonceAccount;
if let Ok(info) = parse_nonce(&account.data) {
match info {
solana_account_decoder::parse_nonce::UiNonceState::Initialized(details) => {
File diff suppressed because it is too large Load Diff
@@ -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<PoolState> {
borsh::from_slice::<PoolState>(&data[..POOL_STATE_SIZE]).ok()
}
pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
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<GlobalConfig> {
pub fn global_config_parser(
account: &AccountPretty,
metadata: EventMetadata,
mut metadata: EventMetadata,
) -> Option<DexEvent> {
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<PlatformConfig> {
pub fn platform_config_parser(
account: &AccountPretty,
metadata: EventMetadata,
mut metadata: EventMetadata,
) -> Option<DexEvent> {
metadata.event_type = EventType::AccountBonkPlatformConfig;
if account.data.len() < PLATFORM_CONFIG_SIZE + 8 {
return None;
}
@@ -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<BondingCurve> {
pub fn bonding_curve_parser(
account: &AccountPretty,
metadata: EventMetadata,
mut metadata: EventMetadata,
) -> Option<DexEvent> {
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<Global> {
borsh::from_slice::<Global>(&data[..GLOBAL_SIZE]).ok()
}
pub fn global_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn global_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
metadata.event_type = EventType::AccountPumpFunGlobal;
if account.data.len() < GLOBAL_SIZE + 8 {
return None;
}
@@ -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<GlobalConfig> {
pub fn global_config_parser(
account: &AccountPretty,
metadata: EventMetadata,
mut metadata: EventMetadata,
) -> Option<DexEvent> {
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<Pool> {
borsh::from_slice::<Pool>(&data[..POOL_SIZE]).ok()
}
pub fn pool_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn pool_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
metadata.event_type = EventType::AccountPumpSwapPool;
if account.data.len() < POOL_SIZE + 8 {
return None;
}
@@ -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<AmmInfo> {
borsh::from_slice::<AmmInfo>(&data[..AMM_INFO_SIZE]).ok()
}
pub fn amm_info_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn amm_info_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
metadata.event_type = EventType::AccountRaydiumAmmV4AmmInfo;
if account.data.len() < AMM_INFO_SIZE {
return None;
}
@@ -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<AmmConfig> {
borsh::from_slice::<AmmConfig>(&data[..AMM_CONFIG_SIZE]).ok()
}
pub fn amm_config_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn amm_config_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
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<PoolState> {
borsh::from_slice::<PoolState>(&data[..POOL_STATE_SIZE]).ok()
}
pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
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<TickArrayState> {
pub fn tick_array_state_parser(
account: &AccountPretty,
metadata: EventMetadata,
mut metadata: EventMetadata,
) -> Option<DexEvent> {
metadata.event_type = EventType::AccountRaydiumClmmTickArrayState;
if account.data.len() < TICK_ARRAY_STATE_SIZE + 8 {
return None;
}
@@ -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<AmmConfig> {
borsh::from_slice::<AmmConfig>(&data[..AMM_CONFIG_SIZE]).ok()
}
pub fn amm_config_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn amm_config_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
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<PoolState> {
borsh::from_slice::<PoolState>(&data[..POOL_STATE_SIZE]).ok()
}
pub fn pool_state_parser(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent> {
pub fn pool_state_parser(account: &AccountPretty, mut metadata: EventMetadata) -> Option<DexEvent> {
metadata.event_type = EventType::AccountRaydiumCpmmPoolState;
if account.data.len() < POOL_STATE_SIZE + 8 {
return None;
}