use crate::streaming::{
event_parser::{
common::{filter::EventTypeFilter, EventMetadata, EventType, ProtocolType},
protocols::{
bonk::parser::BONK_PROGRAM_ID,
pumpfun::parser::PUMPFUN_PROGRAM_ID,
pumpswap::parser::PUMPSWAP_PROGRAM_ID,
raydium_amm_v4::parser::RAYDIUM_AMM_V4_PROGRAM_ID,
raydium_clmm::parser::RAYDIUM_CLMM_PROGRAM_ID,
raydium_cpmm::parser::RAYDIUM_CPMM_PROGRAM_ID,
},
Protocol, DexEvent,
},
grpc::AccountPretty,
};
use solana_sdk::pubkey::Pubkey;
use std::{
collections::HashMap,
sync::{Arc, LazyLock, OnceLock},
};
pub type InnerInstructionEventParser =
fn(data: &[u8], metadata: EventMetadata) -> Option<DexEvent>;
pub type InstructionEventParser =
fn(data: &[u8], accounts: &[Pubkey], metadata: EventMetadata) -> Option<DexEvent>;
#[derive(Debug, Clone)]
pub struct GenericEventParseConfig {
pub program_id: Pubkey,
pub protocol_type: ProtocolType,
pub inner_instruction_discriminator: &'static [u8],
pub instruction_discriminator: &'static [u8],
pub event_type: EventType,
pub inner_instruction_parser: Option<InnerInstructionEventParser>,
pub instruction_parser: Option<InstructionEventParser>,
pub requires_inner_instruction: bool,
}
pub static EVENT_PARSERS: LazyLock<HashMap<Protocol, (Pubkey, &[GenericEventParseConfig])>> =
LazyLock::new(|| {
let mut parsers = HashMap::with_capacity(6);
parsers.insert(
Protocol::PumpSwap,
(
PUMPSWAP_PROGRAM_ID,
crate::streaming::event_parser::protocols::pumpswap::parser::CONFIGS,
),
);
parsers.insert(
Protocol::PumpFun,
(
PUMPFUN_PROGRAM_ID,
crate::streaming::event_parser::protocols::pumpfun::parser::CONFIGS,
),
);
parsers.insert(
Protocol::Bonk,
(BONK_PROGRAM_ID, crate::streaming::event_parser::protocols::bonk::parser::CONFIGS),
);
parsers.insert(
Protocol::RaydiumCpmm,
(
RAYDIUM_CPMM_PROGRAM_ID,
crate::streaming::event_parser::protocols::raydium_cpmm::parser::CONFIGS,
),
);
parsers.insert(
Protocol::RaydiumClmm,
(
RAYDIUM_CLMM_PROGRAM_ID,
crate::streaming::event_parser::protocols::raydium_clmm::parser::CONFIGS,
),
);
parsers.insert(
Protocol::RaydiumAmmV4,
(
RAYDIUM_AMM_V4_PROGRAM_ID,
crate::streaming::event_parser::protocols::raydium_amm_v4::parser::CONFIGS,
),
);
parsers
});
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct CacheKey {
pub protocols: Vec<Protocol>,
pub event_types: Option<Vec<EventType>>,
}
impl CacheKey {
pub fn new(mut protocols: Vec<Protocol>, filter: Option<&EventTypeFilter>) -> Self {
protocols.sort_by_cached_key(|p| format!("{:?}", p));
let event_types = filter.map(|f| {
let mut types = f.include.clone();
types.sort_by_cached_key(|t| format!("{:?}", t));
types
});
Self { protocols, event_types }
}
}
static GLOBAL_PROGRAM_IDS_CACHE: LazyLock<
parking_lot::RwLock<HashMap<CacheKey, Arc<Vec<Pubkey>>>>,
> = LazyLock::new(|| parking_lot::RwLock::new(HashMap::new()));
static GLOBAL_INSTRUCTION_CONFIGS_CACHE: LazyLock<
parking_lot::RwLock<HashMap<CacheKey, Arc<HashMap<Vec<u8>, Vec<GenericEventParseConfig>>>>>,
> = LazyLock::new(|| parking_lot::RwLock::new(HashMap::new()));
pub fn get_global_program_ids(
protocols: &[Protocol],
filter: Option<&EventTypeFilter>,
) -> Arc<Vec<Pubkey>> {
let cache_key = CacheKey::new(protocols.to_vec(), filter);
{
let cache = GLOBAL_PROGRAM_IDS_CACHE.read();
if let Some(program_ids) = cache.get(&cache_key) {
return program_ids.clone();
}
}
let mut program_ids = Vec::with_capacity(protocols.len());
for protocol in protocols {
if let Some(parse) = EVENT_PARSERS.get(protocol) {
program_ids.push(parse.0);
}
}
let program_ids = Arc::new(program_ids);
GLOBAL_PROGRAM_IDS_CACHE.write().insert(cache_key, program_ids.clone());
program_ids
}
pub fn get_global_instruction_configs(
protocols: &[Protocol],
filter: Option<&EventTypeFilter>,
) -> Arc<HashMap<Vec<u8>, Vec<GenericEventParseConfig>>> {
let cache_key = CacheKey::new(protocols.to_vec(), filter);
{
let cache = GLOBAL_INSTRUCTION_CONFIGS_CACHE.read();
if let Some(configs) = cache.get(&cache_key) {
return configs.clone();
}
}
let mut instruction_configs = HashMap::with_capacity(protocols.len() * 4);
for protocol in protocols {
if let Some(parse) = EVENT_PARSERS.get(protocol) {
parse
.1
.iter()
.filter(|config| {
filter.as_ref().map(|f| f.include.contains(&config.event_type)).unwrap_or(true)
})
.for_each(|config| {
instruction_configs
.entry(config.instruction_discriminator.to_vec())
.or_insert_with(Vec::new)
.push(config.clone());
});
}
}
let instruction_configs = Arc::new(instruction_configs);
GLOBAL_INSTRUCTION_CONFIGS_CACHE.write().insert(cache_key, instruction_configs.clone());
instruction_configs
}
#[derive(Debug)]
pub struct AccountPubkeyCache {
cache: Vec<Pubkey>,
}
impl AccountPubkeyCache {
pub fn new() -> Self {
Self {
cache: Vec::with_capacity(32),
}
}
#[inline]
pub fn build_account_pubkeys(
&mut self,
instruction_accounts: &[u8],
all_accounts: &[Pubkey],
) -> &[Pubkey] {
self.cache.clear();
if self.cache.capacity() < instruction_accounts.len() {
self.cache.reserve(instruction_accounts.len() - self.cache.capacity());
}
for &idx in instruction_accounts.iter() {
if (idx as usize) < all_accounts.len() {
self.cache.push(all_accounts[idx as usize]);
}
}
&self.cache
}
}
impl Default for AccountPubkeyCache {
fn default() -> Self {
Self::new()
}
}
thread_local! {
static THREAD_LOCAL_ACCOUNT_CACHE: std::cell::RefCell<AccountPubkeyCache> =
std::cell::RefCell::new(AccountPubkeyCache::new());
}
#[inline]
pub fn build_account_pubkeys_with_cache(
instruction_accounts: &[u8],
all_accounts: &[Pubkey],
) -> Vec<Pubkey> {
THREAD_LOCAL_ACCOUNT_CACHE.with(|cache| {
let mut cache = cache.borrow_mut();
cache.build_account_pubkeys(instruction_accounts, all_accounts).to_vec()
})
}
pub type AccountEventParserFn =
fn(account: &AccountPretty, metadata: EventMetadata) -> Option<DexEvent>;
#[derive(Debug, Clone)]
pub struct AccountEventParseConfig {
pub program_id: Pubkey,
pub protocol_type: ProtocolType,
pub event_type: EventType,
pub account_discriminator: &'static [u8],
pub account_parser: AccountEventParserFn,
}
static PROTOCOL_CONFIGS_CACHE: OnceLock<HashMap<Protocol, Vec<AccountEventParseConfig>>> =
OnceLock::new();
static NONCE_CONFIG: OnceLock<AccountEventParseConfig> = OnceLock::new();
static COMMON_CONFIG: OnceLock<AccountEventParseConfig> = OnceLock::new();
pub fn get_account_configs(
protocols: &[Protocol],
event_type_filter: Option<&EventTypeFilter>,
nonce_parser: AccountEventParserFn,
token_parser: AccountEventParserFn,
) -> Vec<AccountEventParseConfig> {
let protocols_map = PROTOCOL_CONFIGS_CACHE.get_or_init(|| {
let mut map = HashMap::new();
map.insert(Protocol::PumpSwap, vec![
AccountEventParseConfig {
program_id: PUMPSWAP_PROGRAM_ID,
protocol_type: ProtocolType::PumpSwap,
event_type: EventType::AccountPumpSwapGlobalConfig,
account_discriminator: crate::streaming::event_parser::protocols::pumpswap::discriminators::GLOBAL_CONFIG_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::pumpswap::types::global_config_parser,
},
AccountEventParseConfig {
program_id: PUMPSWAP_PROGRAM_ID,
protocol_type: ProtocolType::PumpSwap,
event_type: EventType::AccountPumpSwapPool,
account_discriminator: crate::streaming::event_parser::protocols::pumpswap::discriminators::POOL_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::pumpswap::types::pool_parser,
},
]);
map.insert(Protocol::PumpFun, vec![
AccountEventParseConfig {
program_id: PUMPFUN_PROGRAM_ID,
protocol_type: ProtocolType::PumpFun,
event_type: EventType::AccountPumpFunBondingCurve,
account_discriminator: crate::streaming::event_parser::protocols::pumpfun::discriminators::BONDING_CURVE_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::pumpfun::types::bonding_curve_parser,
},
AccountEventParseConfig {
program_id: PUMPFUN_PROGRAM_ID,
protocol_type: ProtocolType::PumpFun,
event_type: EventType::AccountPumpFunGlobal,
account_discriminator: crate::streaming::event_parser::protocols::pumpfun::discriminators::GLOBAL_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::pumpfun::types::global_parser,
},
]);
map.insert(Protocol::Bonk, vec![
AccountEventParseConfig {
program_id: BONK_PROGRAM_ID,
protocol_type: ProtocolType::Bonk,
event_type: EventType::AccountBonkPoolState,
account_discriminator: crate::streaming::event_parser::protocols::bonk::discriminators::POOL_STATE_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::bonk::types::pool_state_parser,
},
AccountEventParseConfig {
program_id: BONK_PROGRAM_ID,
protocol_type: ProtocolType::Bonk,
event_type: EventType::AccountBonkGlobalConfig,
account_discriminator: crate::streaming::event_parser::protocols::bonk::discriminators::GLOBAL_CONFIG_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::bonk::types::global_config_parser,
},
AccountEventParseConfig {
program_id: BONK_PROGRAM_ID,
protocol_type: ProtocolType::Bonk,
event_type: EventType::AccountBonkPlatformConfig,
account_discriminator: crate::streaming::event_parser::protocols::bonk::discriminators::PLATFORM_CONFIG_ACCOUNT,
account_parser: crate::streaming::event_parser::protocols::bonk::types::platform_config_parser,
},
]);
map.insert(Protocol::RaydiumCpmm, vec![
AccountEventParseConfig {
program_id: RAYDIUM_CPMM_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumCpmm,
event_type: EventType::AccountRaydiumCpmmAmmConfig,
account_discriminator: crate::streaming::event_parser::protocols::raydium_cpmm::discriminators::AMM_CONFIG,
account_parser: crate::streaming::event_parser::protocols::raydium_cpmm::types::amm_config_parser,
},
AccountEventParseConfig {
program_id: RAYDIUM_CPMM_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumCpmm,
event_type: EventType::AccountRaydiumCpmmPoolState,
account_discriminator: crate::streaming::event_parser::protocols::raydium_cpmm::discriminators::POOL_STATE,
account_parser: crate::streaming::event_parser::protocols::raydium_cpmm::types::pool_state_parser,
},
]);
map.insert(Protocol::RaydiumClmm, vec![
AccountEventParseConfig {
program_id: RAYDIUM_CLMM_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumClmm,
event_type: EventType::AccountRaydiumClmmAmmConfig,
account_discriminator: crate::streaming::event_parser::protocols::raydium_clmm::discriminators::AMM_CONFIG,
account_parser: crate::streaming::event_parser::protocols::raydium_clmm::types::amm_config_parser,
},
AccountEventParseConfig {
program_id: RAYDIUM_CLMM_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumClmm,
event_type: EventType::AccountRaydiumClmmPoolState,
account_discriminator: crate::streaming::event_parser::protocols::raydium_clmm::discriminators::POOL_STATE,
account_parser: crate::streaming::event_parser::protocols::raydium_clmm::types::pool_state_parser,
},
AccountEventParseConfig {
program_id: RAYDIUM_CLMM_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumClmm,
event_type: EventType::AccountRaydiumClmmTickArrayState,
account_discriminator: crate::streaming::event_parser::protocols::raydium_clmm::discriminators::TICK_ARRAY_STATE,
account_parser: crate::streaming::event_parser::protocols::raydium_clmm::types::tick_array_state_parser,
},
]);
map.insert(Protocol::RaydiumAmmV4, vec![
AccountEventParseConfig {
program_id: RAYDIUM_AMM_V4_PROGRAM_ID,
protocol_type: ProtocolType::RaydiumAmmV4,
event_type: EventType::AccountRaydiumAmmV4AmmInfo,
account_discriminator: crate::streaming::event_parser::protocols::raydium_amm_v4::discriminators::AMM_INFO,
account_parser: crate::streaming::event_parser::protocols::raydium_amm_v4::types::amm_info_parser,
},
]);
map
});
let mut configs = Vec::new();
let empty_vec = Vec::new();
let estimated_capacity = protocols.len() * 3;
configs.reserve(estimated_capacity);
for protocol in protocols {
let protocol_configs = protocols_map.get(protocol).unwrap_or(&empty_vec);
if event_type_filter.is_none() {
configs.extend(protocol_configs.iter().cloned());
} else {
let filter = event_type_filter.unwrap();
configs.extend(
protocol_configs
.iter()
.filter(|config| filter.include.contains(&config.event_type))
.cloned(),
);
}
}
if event_type_filter.is_none()
|| event_type_filter.unwrap().include.contains(&EventType::NonceAccount)
{
let nonce_config = NONCE_CONFIG.get_or_init(|| AccountEventParseConfig {
program_id: Pubkey::default(),
protocol_type: ProtocolType::Common,
event_type: EventType::NonceAccount,
account_discriminator: &[1, 0, 0, 0, 1, 0, 0, 0],
account_parser: nonce_parser,
});
configs.push(nonce_config.clone());
}
let common_config = COMMON_CONFIG.get_or_init(|| AccountEventParseConfig {
program_id: Pubkey::default(),
protocol_type: ProtocolType::Common,
event_type: EventType::TokenAccount,
account_discriminator: &[],
account_parser: token_parser,
});
configs.push(common_config.clone());
configs
}