Skip to main content

sol_parser_sdk/grpc/
types.rs

1use serde::{Deserialize, Serialize};
2use yellowstone_grpc_proto::geyser::{
3    subscribe_request_filter_accounts_filter::Filter as AccountsFilterOneof,
4    subscribe_request_filter_accounts_filter_memcmp::Data as MemcmpDataOneof,
5    SubscribeRequestFilterAccountsFilter, SubscribeRequestFilterAccountsFilterMemcmp,
6};
7
8/// 事件输出顺序模式
9#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
10pub enum OrderMode {
11    /// 无序模式:收到即输出,超低延迟 (10-20μs)
12    #[default]
13    Unordered,
14    /// 有序模式:按 slot + tx_index 排序后输出
15    /// 同一 slot 内的交易会等待收齐后按 tx_index 排序
16    /// 延迟增加约 1-50ms(取决于 slot 内交易数量)
17    Ordered,
18    /// 流式有序模式:连续序列立即释放,低延迟 + 顺序保证
19    /// 只要收到从 0 开始的连续 tx_index 序列,立即释放
20    /// 延迟约 0.1-5ms,比 Ordered 低 5-50 倍
21    StreamingOrdered,
22    /// 微批次模式:极短时间窗口内收集事件,窗口结束后排序释放
23    /// 窗口大小由 micro_batch_us 配置(默认 100μs)
24    /// 延迟约 50-200μs,接近 Unordered 但保证顺序
25    MicroBatch,
26}
27
28#[derive(Debug, Clone, Serialize, Deserialize)]
29pub struct ClientConfig {
30    /// 是否启用性能监控
31    pub enable_metrics: bool,
32    /// 连接超时时间(毫秒)
33    pub connection_timeout_ms: u64,
34    /// 请求超时时间(毫秒)
35    pub request_timeout_ms: u64,
36    /// 是否启用TLS
37    pub enable_tls: bool,
38    pub max_retries: u32,
39    pub retry_delay_ms: u64,
40    pub max_concurrent_streams: u32,
41    pub keep_alive_interval_ms: u64,
42    pub keep_alive_timeout_ms: u64,
43    pub buffer_size: usize,
44    /// 事件输出顺序模式
45    pub order_mode: OrderMode,
46    /// 有序模式下,slot 超时时间(毫秒)
47    /// 超过此时间未收到新 slot 信号,强制输出当前缓冲的事件
48    pub order_timeout_ms: u64,
49    /// MicroBatch 模式下的时间窗口大小(微秒)
50    /// 默认 100μs,可根据网络状况调整
51    pub micro_batch_us: u64,
52}
53
54impl Default for ClientConfig {
55    fn default() -> Self {
56        Self {
57            enable_metrics: false,
58            connection_timeout_ms: 8000,
59            request_timeout_ms: 15000,
60            enable_tls: true,
61            max_retries: 3,
62            retry_delay_ms: 1000,
63            max_concurrent_streams: 100,
64            keep_alive_interval_ms: 30000,
65            keep_alive_timeout_ms: 5000,
66            buffer_size: 100_000,
67            order_mode: OrderMode::Unordered,
68            order_timeout_ms: 100,
69            micro_batch_us: 100, // 100μs 默认窗口
70        }
71    }
72}
73
74impl ClientConfig {
75    pub fn low_latency() -> Self {
76        Self {
77            enable_metrics: false,
78            connection_timeout_ms: 5000,
79            request_timeout_ms: 10000,
80            enable_tls: true,
81            max_retries: 1,
82            retry_delay_ms: 100,
83            max_concurrent_streams: 200,
84            keep_alive_interval_ms: 10000,
85            keep_alive_timeout_ms: 2000,
86            buffer_size: 100_000,
87            order_mode: OrderMode::Unordered,
88            order_timeout_ms: 50,
89            micro_batch_us: 50, // 50μs 更激进的窗口
90        }
91    }
92
93    pub fn high_throughput() -> Self {
94        Self {
95            enable_metrics: true,
96            connection_timeout_ms: 10000,
97            request_timeout_ms: 30000,
98            enable_tls: true,
99            max_retries: 5,
100            retry_delay_ms: 2000,
101            max_concurrent_streams: 500,
102            keep_alive_interval_ms: 60000,
103            keep_alive_timeout_ms: 10000,
104            buffer_size: 200_000,
105            order_mode: OrderMode::Unordered,
106            order_timeout_ms: 200,
107            micro_batch_us: 200, // 200μs 高吞吐模式
108        }
109    }
110}
111
112#[derive(Debug, Clone)]
113pub struct TransactionFilter {
114    pub account_include: Vec<String>,
115    pub account_exclude: Vec<String>,
116    pub account_required: Vec<String>,
117}
118
119impl TransactionFilter {
120    pub fn new() -> Self {
121        Self {
122            account_include: Vec::new(),
123            account_exclude: Vec::new(),
124            account_required: Vec::new(),
125        }
126    }
127
128    pub fn include_account(mut self, account: impl Into<String>) -> Self {
129        self.account_include.push(account.into());
130        self
131    }
132
133    pub fn exclude_account(mut self, account: impl Into<String>) -> Self {
134        self.account_exclude.push(account.into());
135        self
136    }
137
138    pub fn require_account(mut self, account: impl Into<String>) -> Self {
139        self.account_required.push(account.into());
140        self
141    }
142
143    /// 从程序ID列表创建过滤器
144    pub fn from_program_ids(program_ids: Vec<String>) -> Self {
145        Self {
146            account_include: program_ids,
147            account_exclude: Vec::new(),
148            account_required: Vec::new(),
149        }
150    }
151}
152
153impl Default for TransactionFilter {
154    fn default() -> Self {
155        Self::new()
156    }
157}
158
159#[derive(Debug, Clone)]
160pub struct AccountFilter {
161    pub account: Vec<String>,
162    pub owner: Vec<String>,
163    pub filters: Vec<SubscribeRequestFilterAccountsFilter>,
164}
165
166impl AccountFilter {
167    pub fn new() -> Self {
168        Self { account: Vec::new(), owner: Vec::new(), filters: Vec::new() }
169    }
170
171    pub fn add_account(mut self, account: impl Into<String>) -> Self {
172        self.account.push(account.into());
173        self
174    }
175
176    pub fn add_owner(mut self, owner: impl Into<String>) -> Self {
177        self.owner.push(owner.into());
178        self
179    }
180
181    pub fn add_filter(mut self, filter: SubscribeRequestFilterAccountsFilter) -> Self {
182        self.filters.push(filter);
183        self
184    }
185
186    /// 从程序ID列表创建所有者过滤器
187    pub fn from_program_owners(program_ids: Vec<String>) -> Self {
188        Self { account: Vec::new(), owner: program_ids, filters: Vec::new() }
189    }
190}
191
192impl Default for AccountFilter {
193    fn default() -> Self {
194        Self::new()
195    }
196}
197
198/// Build a memcmp account filter for use in `AccountFilter::filters`.
199/// ATA accounts have mint at offset 0; PumpSwap pool accounts often use offset 32 for mint/pubkey.
200#[inline]
201pub fn account_filter_memcmp(offset: u64, bytes: Vec<u8>) -> SubscribeRequestFilterAccountsFilter {
202    SubscribeRequestFilterAccountsFilter {
203        filter: Some(AccountsFilterOneof::Memcmp(SubscribeRequestFilterAccountsFilterMemcmp {
204            offset,
205            data: Some(MemcmpDataOneof::Bytes(bytes)),
206        })),
207    }
208}
209
210#[derive(Debug, Clone)]
211pub struct AccountFilterData {
212    pub memcmp: Option<AccountFilterMemcmp>,
213    pub datasize: Option<u64>,
214}
215
216#[derive(Debug, Clone)]
217pub struct AccountFilterMemcmp {
218    pub offset: u64,
219    pub bytes: Vec<u8>,
220}
221
222#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
223pub enum Protocol {
224    PumpFun,
225    PumpSwap,
226    PumpFees,
227    RaydiumLaunchlab,
228    RaydiumCpmm,
229    RaydiumClmm,
230    RaydiumAmmV4,
231    OrcaWhirlpool,
232    MeteoraPools,
233    MeteoraDammV2,
234    MeteoraDlmm,
235    MeteoraDbc,
236}
237
238#[derive(Debug, Clone, Copy, PartialEq, Eq)]
239#[non_exhaustive]
240pub enum EventType {
241    // Block events
242    BlockMeta,
243
244    // RaydiumLaunchlab events
245    RaydiumLaunchlabTrade,
246    RaydiumLaunchlabPoolCreate,
247    RaydiumLaunchlabMigrateAmm,
248
249    // PumpFun events
250    PumpFunTrade,         // All trade events (backward compatible)
251    PumpFunBuy,           // Buy events only (filter by ix_name)
252    PumpFunSell,          // Sell events only (filter by ix_name)
253    PumpFunBuyExactSolIn, // BuyExactSolIn events only (filter by ix_name)
254    PumpFunCreate,
255    PumpFunCreateV2, // SPL-22 / Mayhem create
256    PumpFunComplete,
257    PumpFunMigrate,
258    /// Pump fees(`pfeeUx...`,`idls/pump_fees.json` Program data events)
259    PumpFeesCreateFeeSharingConfig,
260    PumpFeesInitializeFeeConfig,
261    PumpFeesResetFeeSharingConfig,
262    PumpFeesRevokeFeeSharingAuthority,
263    PumpFeesTransferFeeSharingAuthority,
264    PumpFeesUpdateAdmin,
265    PumpFeesUpdateFeeConfig,
266    PumpFeesUpdateFeeShares,
267    PumpFeesUpsertFeeTiers,
268    /// Pump.fun:`migrateBondingCurveCreatorEvent`
269    PumpFunMigrateBondingCurveCreator,
270
271    // PumpSwap events
272    PumpSwapTrade,
273    PumpSwapBuy,
274    PumpSwapSell,
275    PumpSwapCreatePool,
276    PumpSwapLiquidityAdded,
277    PumpSwapLiquidityRemoved,
278    // PumpSwapPoolUpdated,
279    // PumpSwapFeesClaimed,
280
281    // Raydium CPMM events
282    RaydiumCpmmSwap,
283    RaydiumCpmmDeposit,
284    RaydiumCpmmWithdraw,
285    RaydiumCpmmInitialize,
286
287    // Raydium CLMM events
288    RaydiumClmmSwap,
289    RaydiumClmmCreatePool,
290    RaydiumClmmOpenPosition,
291    RaydiumClmmClosePosition,
292    RaydiumClmmIncreaseLiquidity,
293    RaydiumClmmDecreaseLiquidity,
294    RaydiumClmmLiquidityChange,
295    RaydiumClmmConfigChange,
296    RaydiumClmmCreatePersonalPosition,
297    RaydiumClmmLiquidityCalculate,
298    RaydiumClmmOpenLimitOrder,
299    RaydiumClmmIncreaseLimitOrder,
300    RaydiumClmmDecreaseLimitOrder,
301    RaydiumClmmSettleLimitOrder,
302    RaydiumClmmUpdateRewardInfos,
303    RaydiumClmmOpenPositionWithTokenExtNft,
304    RaydiumClmmCollectFee,
305
306    // Raydium AMM V4 events
307    RaydiumAmmV4Swap,
308    RaydiumAmmV4Deposit,
309    RaydiumAmmV4Withdraw,
310    RaydiumAmmV4Initialize2,
311    RaydiumAmmV4WithdrawPnl,
312
313    // Orca Whirlpool events
314    OrcaWhirlpoolSwap,
315    OrcaWhirlpoolLiquidityIncreased,
316    OrcaWhirlpoolLiquidityDecreased,
317    OrcaWhirlpoolPoolInitialized,
318
319    // Meteora events
320    MeteoraPoolsSwap,
321    MeteoraPoolsAddLiquidity,
322    MeteoraPoolsRemoveLiquidity,
323    MeteoraPoolsBootstrapLiquidity,
324    MeteoraPoolsPoolCreated,
325    MeteoraPoolsSetPoolFees,
326
327    // Meteora DAMM V2 events
328    MeteoraDammV2Swap,
329    MeteoraDammV2AddLiquidity,
330    MeteoraDammV2RemoveLiquidity,
331    MeteoraDammV2InitializePool,
332    MeteoraDammV2CreatePosition,
333    MeteoraDammV2ClosePosition,
334    MeteoraDammV2UpdateDelegatePermission,
335    MeteoraDammV2WithdrawDeadLiquidityReward,
336    MeteoraDammV2CreateConfig,
337    MeteoraDammV2CreateDynamicConfig,
338    // MeteoraDammV2ClaimPositionFee,
339    // MeteoraDammV2InitializeReward,
340    // MeteoraDammV2FundReward,
341    // MeteoraDammV2ClaimReward,
342
343    // Meteora DBC events
344    MeteoraDbcSwap,
345    MeteoraDbcInitializePool,
346    MeteoraDbcCurveComplete,
347
348    // Meteora DLMM events
349    MeteoraDlmmSwap,
350    MeteoraDlmmAddLiquidity,
351    MeteoraDlmmRemoveLiquidity,
352    MeteoraDlmmInitializePool,
353    MeteoraDlmmInitializeBinArray,
354    MeteoraDlmmCreatePosition,
355    MeteoraDlmmClosePosition,
356    MeteoraDlmmClaimFee,
357
358    // Account events
359    TokenAccount,
360    TokenInfo,
361    NonceAccount,
362    AccountPumpFunGlobal,
363    AccountPumpFunBondingCurve,
364    AccountPumpFunFeeConfig,
365    AccountPumpFunSharingConfig,
366    AccountPumpFunGlobalVolumeAccumulator,
367    AccountPumpFunUserVolumeAccumulator,
368
369    AccountPumpSwapGlobalConfig,
370    AccountPumpSwapPool,
371    AccountRaydiumClmmAmmConfig,
372    AccountRaydiumClmmPoolState,
373    AccountRaydiumClmmTickArrayState,
374    AccountRaydiumCpmmAmmConfig,
375    AccountRaydiumCpmmPoolState,
376    AccountOrcaWhirlpool,
377    AccountOrcaPosition,
378    AccountOrcaTickArray,
379    AccountOrcaFeeTier,
380    AccountOrcaWhirlpoolsConfig,
381}
382
383#[derive(Debug, Clone)]
384pub struct EventTypeFilter {
385    pub include_only: Option<Vec<EventType>>,
386    pub exclude_types: Option<Vec<EventType>>,
387}
388
389impl EventTypeFilter {
390    pub fn include_only(types: Vec<EventType>) -> Self {
391        Self { include_only: Some(types), exclude_types: None }
392    }
393
394    pub fn exclude_types(types: Vec<EventType>) -> Self {
395        Self { include_only: None, exclude_types: Some(types) }
396    }
397
398    #[inline]
399    fn includes_any(&self, event_types: &[EventType]) -> bool {
400        event_types.iter().any(|event_type| self.should_include(*event_type))
401    }
402
403    pub fn should_include(&self, event_type: EventType) -> bool {
404        if let Some(ref include_only) = self.include_only {
405            // Direct match
406            if include_only.contains(&event_type) {
407                return true;
408            }
409            if matches!(
410                event_type,
411                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
412            ) {
413                if pumpfun_trade_filter_is_generic(include_only) {
414                    return true;
415                }
416                if event_type == EventType::PumpFunBuyExactSolIn
417                    && pumpfun_buy_filter_is_generic(include_only)
418                {
419                    return true;
420                }
421                return false;
422            }
423            if is_pumpfun_create_family(event_type) {
424                return include_only.iter().any(|t| is_pumpfun_create_family(*t));
425            }
426            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell) {
427                return include_only.contains(&EventType::PumpSwapTrade);
428            }
429            return false;
430        }
431
432        if let Some(ref exclude_types) = self.exclude_types {
433            if exclude_types.contains(&event_type) {
434                return false;
435            }
436            if matches!(
437                event_type,
438                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
439            ) && exclude_types.contains(&EventType::PumpFunTrade)
440            {
441                return false;
442            }
443            if event_type == EventType::PumpFunBuyExactSolIn
444                && exclude_types.contains(&EventType::PumpFunBuy)
445            {
446                return false;
447            }
448            if is_pumpfun_create_family(event_type)
449                && exclude_types.iter().any(|t| is_pumpfun_create_family(*t))
450            {
451                return false;
452            }
453            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell)
454                && exclude_types.contains(&EventType::PumpSwapTrade)
455            {
456                return false;
457            }
458            return true;
459        }
460
461        true
462    }
463
464    pub fn should_include_dex_event(&self, event: &crate::core::events::DexEvent) -> bool {
465        let Some(event_type) = event_type_from_dex_event(event) else { return true };
466        self.should_include(event_type)
467    }
468
469    #[inline]
470    pub fn includes_block_meta(&self) -> bool {
471        if let Some(ref include_only) = self.include_only {
472            return include_only.contains(&EventType::BlockMeta);
473        }
474        false
475    }
476
477    #[inline]
478    pub fn normalize_dex_event(
479        &self,
480        event: crate::core::events::DexEvent,
481    ) -> crate::core::events::DexEvent {
482        use crate::core::events::DexEvent;
483
484        let Some(ref include_only) = self.include_only else { return event };
485        if pumpfun_trade_filter_is_generic(include_only) {
486            return match event {
487                DexEvent::PumpFunBuy(t)
488                | DexEvent::PumpFunSell(t)
489                | DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunTrade(t),
490                other => other,
491            };
492        }
493        if pumpfun_buy_filter_is_generic(include_only) {
494            return match event {
495                DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunBuy(t),
496                other => other,
497            };
498        }
499
500        event
501    }
502
503    #[inline]
504    pub fn includes_pumpfun(&self) -> bool {
505        self.includes_any(&[
506            EventType::PumpFunTrade,
507            EventType::PumpFunBuy,
508            EventType::PumpFunSell,
509            EventType::PumpFunBuyExactSolIn,
510            EventType::PumpFunCreate,
511            EventType::PumpFunCreateV2,
512            EventType::PumpFunComplete,
513            EventType::PumpFunMigrate,
514            EventType::PumpFunMigrateBondingCurveCreator,
515        ])
516    }
517
518    #[inline]
519    pub fn includes_meteora_damm_v2(&self) -> bool {
520        self.includes_any(&[
521            EventType::MeteoraDammV2Swap,
522            EventType::MeteoraDammV2AddLiquidity,
523            EventType::MeteoraDammV2CreatePosition,
524            EventType::MeteoraDammV2ClosePosition,
525            EventType::MeteoraDammV2InitializePool,
526            EventType::MeteoraDammV2RemoveLiquidity,
527            EventType::MeteoraDammV2UpdateDelegatePermission,
528            EventType::MeteoraDammV2WithdrawDeadLiquidityReward,
529            EventType::MeteoraDammV2CreateConfig,
530            EventType::MeteoraDammV2CreateDynamicConfig,
531        ])
532    }
533
534    #[inline]
535    pub fn includes_pump_fees(&self) -> bool {
536        self.includes_any(&[
537            EventType::PumpFeesCreateFeeSharingConfig,
538            EventType::PumpFeesInitializeFeeConfig,
539            EventType::PumpFeesResetFeeSharingConfig,
540            EventType::PumpFeesRevokeFeeSharingAuthority,
541            EventType::PumpFeesTransferFeeSharingAuthority,
542            EventType::PumpFeesUpdateAdmin,
543            EventType::PumpFeesUpdateFeeConfig,
544            EventType::PumpFeesUpdateFeeShares,
545            EventType::PumpFeesUpsertFeeTiers,
546        ])
547    }
548
549    /// Check if PumpSwap protocol events are included in the filter
550    #[inline]
551    pub fn includes_pumpswap(&self) -> bool {
552        self.includes_any(&[
553            EventType::PumpSwapTrade,
554            EventType::PumpSwapBuy,
555            EventType::PumpSwapSell,
556            EventType::PumpSwapCreatePool,
557            EventType::PumpSwapLiquidityAdded,
558            EventType::PumpSwapLiquidityRemoved,
559        ])
560    }
561
562    /// Check if Raydium LaunchLab events are included in the filter.
563    #[inline]
564    pub fn includes_raydium_launchlab(&self) -> bool {
565        self.includes_any(&[
566            EventType::RaydiumLaunchlabTrade,
567            EventType::RaydiumLaunchlabPoolCreate,
568            EventType::RaydiumLaunchlabMigrateAmm,
569        ])
570    }
571
572    #[inline]
573    pub fn includes_raydium_cpmm(&self) -> bool {
574        self.includes_any(&[
575            EventType::RaydiumCpmmSwap,
576            EventType::RaydiumCpmmDeposit,
577            EventType::RaydiumCpmmWithdraw,
578            EventType::RaydiumCpmmInitialize,
579        ])
580    }
581
582    #[inline]
583    pub fn includes_raydium_clmm(&self) -> bool {
584        self.includes_any(&[
585            EventType::RaydiumClmmSwap,
586            EventType::RaydiumClmmCreatePool,
587            EventType::RaydiumClmmOpenPosition,
588            EventType::RaydiumClmmClosePosition,
589            EventType::RaydiumClmmIncreaseLiquidity,
590            EventType::RaydiumClmmDecreaseLiquidity,
591            EventType::RaydiumClmmLiquidityChange,
592            EventType::RaydiumClmmConfigChange,
593            EventType::RaydiumClmmCreatePersonalPosition,
594            EventType::RaydiumClmmLiquidityCalculate,
595            EventType::RaydiumClmmOpenLimitOrder,
596            EventType::RaydiumClmmIncreaseLimitOrder,
597            EventType::RaydiumClmmDecreaseLimitOrder,
598            EventType::RaydiumClmmSettleLimitOrder,
599            EventType::RaydiumClmmUpdateRewardInfos,
600            EventType::RaydiumClmmOpenPositionWithTokenExtNft,
601            EventType::RaydiumClmmCollectFee,
602        ])
603    }
604
605    #[inline]
606    pub fn includes_raydium_amm_v4(&self) -> bool {
607        self.includes_any(&[
608            EventType::RaydiumAmmV4Swap,
609            EventType::RaydiumAmmV4Deposit,
610            EventType::RaydiumAmmV4Withdraw,
611            EventType::RaydiumAmmV4Initialize2,
612            EventType::RaydiumAmmV4WithdrawPnl,
613        ])
614    }
615
616    #[inline]
617    pub fn includes_orca_whirlpool(&self) -> bool {
618        self.includes_any(&[
619            EventType::OrcaWhirlpoolSwap,
620            EventType::OrcaWhirlpoolLiquidityIncreased,
621            EventType::OrcaWhirlpoolLiquidityDecreased,
622            EventType::OrcaWhirlpoolPoolInitialized,
623        ])
624    }
625
626    #[inline]
627    pub fn includes_meteora_pools(&self) -> bool {
628        self.includes_any(&[
629            EventType::MeteoraPoolsSwap,
630            EventType::MeteoraPoolsAddLiquidity,
631            EventType::MeteoraPoolsRemoveLiquidity,
632            EventType::MeteoraPoolsBootstrapLiquidity,
633            EventType::MeteoraPoolsPoolCreated,
634            EventType::MeteoraPoolsSetPoolFees,
635        ])
636    }
637
638    #[inline]
639    pub fn includes_meteora_dlmm(&self) -> bool {
640        self.includes_any(&[
641            EventType::MeteoraDlmmSwap,
642            EventType::MeteoraDlmmAddLiquidity,
643            EventType::MeteoraDlmmRemoveLiquidity,
644            EventType::MeteoraDlmmInitializePool,
645            EventType::MeteoraDlmmInitializeBinArray,
646            EventType::MeteoraDlmmCreatePosition,
647            EventType::MeteoraDlmmClosePosition,
648            EventType::MeteoraDlmmClaimFee,
649        ])
650    }
651
652    #[inline]
653    pub fn includes_meteora_dbc(&self) -> bool {
654        self.includes_any(&[
655            EventType::MeteoraDbcSwap,
656            EventType::MeteoraDbcInitializePool,
657            EventType::MeteoraDbcCurveComplete,
658        ])
659    }
660}
661
662#[inline]
663fn pumpfun_trade_filter_is_generic(include_only: &[EventType]) -> bool {
664    include_only.contains(&EventType::PumpFunTrade)
665        && !include_only.iter().any(|t| {
666            matches!(
667                t,
668                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
669            )
670        })
671}
672
673#[inline]
674fn pumpfun_buy_filter_is_generic(include_only: &[EventType]) -> bool {
675    include_only.contains(&EventType::PumpFunBuy)
676        && !include_only.contains(&EventType::PumpFunBuyExactSolIn)
677}
678
679#[inline]
680fn is_pumpfun_create_family(event_type: EventType) -> bool {
681    matches!(event_type, EventType::PumpFunCreate | EventType::PumpFunCreateV2)
682}
683
684#[inline]
685pub fn event_type_from_dex_event(event: &crate::core::events::DexEvent) -> Option<EventType> {
686    use crate::core::events::DexEvent;
687    match event {
688        DexEvent::PumpFunCreate(_) => Some(EventType::PumpFunCreate),
689        DexEvent::PumpFunCreateV2(_) => Some(EventType::PumpFunCreateV2),
690        DexEvent::PumpFunTrade(_) => Some(EventType::PumpFunTrade),
691        DexEvent::PumpFunBuy(_) => Some(EventType::PumpFunBuy),
692        DexEvent::PumpFunSell(_) => Some(EventType::PumpFunSell),
693        DexEvent::PumpFunBuyExactSolIn(_) => Some(EventType::PumpFunBuyExactSolIn),
694        DexEvent::PumpFunMigrate(_) => Some(EventType::PumpFunMigrate),
695        DexEvent::PumpFeesCreateFeeSharingConfig(_) => {
696            Some(EventType::PumpFeesCreateFeeSharingConfig)
697        }
698        DexEvent::PumpFeesInitializeFeeConfig(_) => Some(EventType::PumpFeesInitializeFeeConfig),
699        DexEvent::PumpFeesResetFeeSharingConfig(_) => {
700            Some(EventType::PumpFeesResetFeeSharingConfig)
701        }
702        DexEvent::PumpFeesRevokeFeeSharingAuthority(_) => {
703            Some(EventType::PumpFeesRevokeFeeSharingAuthority)
704        }
705        DexEvent::PumpFeesTransferFeeSharingAuthority(_) => {
706            Some(EventType::PumpFeesTransferFeeSharingAuthority)
707        }
708        DexEvent::PumpFeesUpdateAdmin(_) => Some(EventType::PumpFeesUpdateAdmin),
709        DexEvent::PumpFeesUpdateFeeConfig(_) => Some(EventType::PumpFeesUpdateFeeConfig),
710        DexEvent::PumpFeesUpdateFeeShares(_) => Some(EventType::PumpFeesUpdateFeeShares),
711        DexEvent::PumpFeesUpsertFeeTiers(_) => Some(EventType::PumpFeesUpsertFeeTiers),
712        DexEvent::PumpFunMigrateBondingCurveCreator(_) => {
713            Some(EventType::PumpFunMigrateBondingCurveCreator)
714        }
715        DexEvent::PumpFunGlobalAccount(_) => Some(EventType::AccountPumpFunGlobal),
716        DexEvent::PumpFunBondingCurveAccount(_) => Some(EventType::AccountPumpFunBondingCurve),
717        DexEvent::PumpFunFeeConfigAccount(_) => Some(EventType::AccountPumpFunFeeConfig),
718        DexEvent::PumpFunSharingConfigAccount(_) => Some(EventType::AccountPumpFunSharingConfig),
719        DexEvent::PumpFunGlobalVolumeAccumulatorAccount(_) => {
720            Some(EventType::AccountPumpFunGlobalVolumeAccumulator)
721        }
722        DexEvent::PumpFunUserVolumeAccumulatorAccount(_) => {
723            Some(EventType::AccountPumpFunUserVolumeAccumulator)
724        }
725        DexEvent::PumpSwapTrade(_) => Some(EventType::PumpSwapTrade),
726        DexEvent::PumpSwapBuy(_) => Some(EventType::PumpSwapBuy),
727        DexEvent::PumpSwapSell(_) => Some(EventType::PumpSwapSell),
728        DexEvent::PumpSwapCreatePool(_) => Some(EventType::PumpSwapCreatePool),
729        DexEvent::PumpSwapLiquidityAdded(_) => Some(EventType::PumpSwapLiquidityAdded),
730        DexEvent::PumpSwapLiquidityRemoved(_) => Some(EventType::PumpSwapLiquidityRemoved),
731        DexEvent::MeteoraDammV2Swap(_) => Some(EventType::MeteoraDammV2Swap),
732        DexEvent::MeteoraDammV2CreatePosition(_) => Some(EventType::MeteoraDammV2CreatePosition),
733        DexEvent::MeteoraDammV2ClosePosition(_) => Some(EventType::MeteoraDammV2ClosePosition),
734        DexEvent::MeteoraDammV2AddLiquidity(_) => Some(EventType::MeteoraDammV2AddLiquidity),
735        DexEvent::MeteoraDammV2RemoveLiquidity(_) => Some(EventType::MeteoraDammV2RemoveLiquidity),
736        DexEvent::MeteoraDammV2InitializePool(_) => Some(EventType::MeteoraDammV2InitializePool),
737        DexEvent::MeteoraDammV2UpdateDelegatePermission(_) => {
738            Some(EventType::MeteoraDammV2UpdateDelegatePermission)
739        }
740        DexEvent::MeteoraDammV2WithdrawDeadLiquidityReward(_) => {
741            Some(EventType::MeteoraDammV2WithdrawDeadLiquidityReward)
742        }
743        DexEvent::MeteoraDammV2CreateConfig(_) => Some(EventType::MeteoraDammV2CreateConfig),
744        DexEvent::MeteoraDammV2CreateDynamicConfig(_) => {
745            Some(EventType::MeteoraDammV2CreateDynamicConfig)
746        }
747        DexEvent::MeteoraDbcSwap(_) => Some(EventType::MeteoraDbcSwap),
748        DexEvent::MeteoraDbcInitializePool(_) => Some(EventType::MeteoraDbcInitializePool),
749        DexEvent::MeteoraDbcCurveComplete(_) => Some(EventType::MeteoraDbcCurveComplete),
750        DexEvent::RaydiumLaunchlabTrade(_) => Some(EventType::RaydiumLaunchlabTrade),
751        DexEvent::RaydiumLaunchlabPoolCreate(_) => Some(EventType::RaydiumLaunchlabPoolCreate),
752        DexEvent::RaydiumLaunchlabMigrateAmm(_) => Some(EventType::RaydiumLaunchlabMigrateAmm),
753        DexEvent::RaydiumClmmSwap(_) => Some(EventType::RaydiumClmmSwap),
754        DexEvent::RaydiumClmmCreatePool(_) => Some(EventType::RaydiumClmmCreatePool),
755        DexEvent::RaydiumClmmOpenPosition(_) => Some(EventType::RaydiumClmmOpenPosition),
756        DexEvent::RaydiumClmmOpenPositionWithTokenExtNft(_) => {
757            Some(EventType::RaydiumClmmOpenPositionWithTokenExtNft)
758        }
759        DexEvent::RaydiumClmmClosePosition(_) => Some(EventType::RaydiumClmmClosePosition),
760        DexEvent::RaydiumClmmIncreaseLiquidity(_) => Some(EventType::RaydiumClmmIncreaseLiquidity),
761        DexEvent::RaydiumClmmDecreaseLiquidity(_) => Some(EventType::RaydiumClmmDecreaseLiquidity),
762        DexEvent::RaydiumClmmLiquidityChange(_) => Some(EventType::RaydiumClmmLiquidityChange),
763        DexEvent::RaydiumClmmConfigChange(_) => Some(EventType::RaydiumClmmConfigChange),
764        DexEvent::RaydiumClmmCreatePersonalPosition(_) => {
765            Some(EventType::RaydiumClmmCreatePersonalPosition)
766        }
767        DexEvent::RaydiumClmmLiquidityCalculate(_) => {
768            Some(EventType::RaydiumClmmLiquidityCalculate)
769        }
770        DexEvent::RaydiumClmmOpenLimitOrder(_) => Some(EventType::RaydiumClmmOpenLimitOrder),
771        DexEvent::RaydiumClmmIncreaseLimitOrder(_) => {
772            Some(EventType::RaydiumClmmIncreaseLimitOrder)
773        }
774        DexEvent::RaydiumClmmDecreaseLimitOrder(_) => {
775            Some(EventType::RaydiumClmmDecreaseLimitOrder)
776        }
777        DexEvent::RaydiumClmmSettleLimitOrder(_) => Some(EventType::RaydiumClmmSettleLimitOrder),
778        DexEvent::RaydiumClmmUpdateRewardInfos(_) => Some(EventType::RaydiumClmmUpdateRewardInfos),
779        DexEvent::RaydiumClmmCollectFee(_) => Some(EventType::RaydiumClmmCollectFee),
780        DexEvent::RaydiumClmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumClmmAmmConfig),
781        DexEvent::RaydiumClmmPoolStateAccount(_) => Some(EventType::AccountRaydiumClmmPoolState),
782        DexEvent::RaydiumClmmTickArrayStateAccount(_) => {
783            Some(EventType::AccountRaydiumClmmTickArrayState)
784        }
785        DexEvent::RaydiumCpmmSwap(_) => Some(EventType::RaydiumCpmmSwap),
786        DexEvent::RaydiumCpmmDeposit(_) => Some(EventType::RaydiumCpmmDeposit),
787        DexEvent::RaydiumCpmmWithdraw(_) => Some(EventType::RaydiumCpmmWithdraw),
788        DexEvent::RaydiumCpmmInitialize(_) => Some(EventType::RaydiumCpmmInitialize),
789        DexEvent::RaydiumCpmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumCpmmAmmConfig),
790        DexEvent::RaydiumCpmmPoolStateAccount(_) => Some(EventType::AccountRaydiumCpmmPoolState),
791        DexEvent::RaydiumAmmV4Swap(_) => Some(EventType::RaydiumAmmV4Swap),
792        DexEvent::RaydiumAmmV4Deposit(_) => Some(EventType::RaydiumAmmV4Deposit),
793        DexEvent::RaydiumAmmV4Initialize2(_) => Some(EventType::RaydiumAmmV4Initialize2),
794        DexEvent::RaydiumAmmV4Withdraw(_) => Some(EventType::RaydiumAmmV4Withdraw),
795        DexEvent::RaydiumAmmV4WithdrawPnl(_) => Some(EventType::RaydiumAmmV4WithdrawPnl),
796        DexEvent::OrcaWhirlpoolSwap(_) => Some(EventType::OrcaWhirlpoolSwap),
797        DexEvent::OrcaWhirlpoolLiquidityIncreased(_) => {
798            Some(EventType::OrcaWhirlpoolLiquidityIncreased)
799        }
800        DexEvent::OrcaWhirlpoolLiquidityDecreased(_) => {
801            Some(EventType::OrcaWhirlpoolLiquidityDecreased)
802        }
803        DexEvent::OrcaWhirlpoolPoolInitialized(_) => Some(EventType::OrcaWhirlpoolPoolInitialized),
804        DexEvent::OrcaWhirlpoolAccount(_) => Some(EventType::AccountOrcaWhirlpool),
805        DexEvent::OrcaPositionAccount(_) => Some(EventType::AccountOrcaPosition),
806        DexEvent::OrcaTickArrayAccount(_) => Some(EventType::AccountOrcaTickArray),
807        DexEvent::OrcaFeeTierAccount(_) => Some(EventType::AccountOrcaFeeTier),
808        DexEvent::OrcaWhirlpoolsConfigAccount(_) => Some(EventType::AccountOrcaWhirlpoolsConfig),
809        DexEvent::MeteoraPoolsSwap(_) => Some(EventType::MeteoraPoolsSwap),
810        DexEvent::MeteoraPoolsAddLiquidity(_) => Some(EventType::MeteoraPoolsAddLiquidity),
811        DexEvent::MeteoraPoolsRemoveLiquidity(_) => Some(EventType::MeteoraPoolsRemoveLiquidity),
812        DexEvent::MeteoraPoolsBootstrapLiquidity(_) => {
813            Some(EventType::MeteoraPoolsBootstrapLiquidity)
814        }
815        DexEvent::MeteoraPoolsPoolCreated(_) => Some(EventType::MeteoraPoolsPoolCreated),
816        DexEvent::MeteoraPoolsSetPoolFees(_) => Some(EventType::MeteoraPoolsSetPoolFees),
817        DexEvent::MeteoraDlmmSwap(_) => Some(EventType::MeteoraDlmmSwap),
818        DexEvent::MeteoraDlmmAddLiquidity(_) => Some(EventType::MeteoraDlmmAddLiquidity),
819        DexEvent::MeteoraDlmmRemoveLiquidity(_) => Some(EventType::MeteoraDlmmRemoveLiquidity),
820        DexEvent::MeteoraDlmmInitializePool(_) => Some(EventType::MeteoraDlmmInitializePool),
821        DexEvent::MeteoraDlmmInitializeBinArray(_) => {
822            Some(EventType::MeteoraDlmmInitializeBinArray)
823        }
824        DexEvent::MeteoraDlmmCreatePosition(_) => Some(EventType::MeteoraDlmmCreatePosition),
825        DexEvent::MeteoraDlmmClosePosition(_) => Some(EventType::MeteoraDlmmClosePosition),
826        DexEvent::MeteoraDlmmClaimFee(_) => Some(EventType::MeteoraDlmmClaimFee),
827        DexEvent::TokenAccount(_) => Some(EventType::TokenAccount),
828        DexEvent::TokenInfo(_) => Some(EventType::TokenInfo),
829        DexEvent::NonceAccount(_) => Some(EventType::NonceAccount),
830        DexEvent::PumpSwapGlobalConfigAccount(_) => Some(EventType::AccountPumpSwapGlobalConfig),
831        DexEvent::PumpSwapPoolAccount(_) => Some(EventType::AccountPumpSwapPool),
832        DexEvent::BlockMeta(_) => Some(EventType::BlockMeta),
833        DexEvent::Error(_) => None,
834    }
835}
836
837#[cfg(test)]
838mod event_type_filter_tests {
839    use super::*;
840
841    #[test]
842    fn generic_trade_filters_cover_specific_trade_variants() {
843        let pump = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
844        assert!(pump.should_include(EventType::PumpFunTrade));
845        assert!(pump.should_include(EventType::PumpFunBuy));
846        assert!(pump.should_include(EventType::PumpFunSell));
847        assert!(pump.should_include(EventType::PumpFunBuyExactSolIn));
848
849        let pump_specific = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
850        assert!(!pump_specific.should_include(EventType::PumpFunTrade));
851        assert!(pump_specific.should_include(EventType::PumpFunBuy));
852        assert!(!pump_specific.should_include(EventType::PumpFunSell));
853        assert!(pump_specific.should_include(EventType::PumpFunBuyExactSolIn));
854
855        let pump_exact_buy = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
856        assert!(!pump_exact_buy.should_include(EventType::PumpFunTrade));
857        assert!(!pump_exact_buy.should_include(EventType::PumpFunBuy));
858        assert!(!pump_exact_buy.should_include(EventType::PumpFunSell));
859        assert!(pump_exact_buy.should_include(EventType::PumpFunBuyExactSolIn));
860
861        let pumpswap = EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]);
862        assert!(pumpswap.should_include(EventType::PumpSwapBuy));
863        assert!(pumpswap.should_include(EventType::PumpSwapSell));
864
865        let exclude_pumpswap = EventTypeFilter::exclude_types(vec![EventType::PumpSwapTrade]);
866        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapBuy));
867        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapSell));
868    }
869
870    #[test]
871    fn generic_pumpfun_trade_filter_normalizes_specific_variants() {
872        use crate::core::events::{DexEvent, PumpFunTradeEvent};
873
874        let filter = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
875        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
876        assert!(matches!(filter.normalize_dex_event(event), DexEvent::PumpFunTrade(_)));
877
878        let specific_filter =
879            EventTypeFilter::include_only(vec![EventType::PumpFunTrade, EventType::PumpFunBuy]);
880        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
881        assert!(matches!(specific_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
882
883        let buy_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
884        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
885            is_buy: true,
886            ..Default::default()
887        });
888        assert!(matches!(buy_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
889
890        let exact_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
891        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
892            is_buy: true,
893            ..Default::default()
894        });
895        assert!(matches!(
896            exact_filter.normalize_dex_event(event),
897            DexEvent::PumpFunBuyExactSolIn(_)
898        ));
899
900        let create_and_trade_filter =
901            EventTypeFilter::include_only(vec![EventType::PumpFunCreate, EventType::PumpFunTrade]);
902        let event =
903            DexEvent::PumpFunSell(PumpFunTradeEvent { is_buy: false, ..Default::default() });
904        assert!(matches!(
905            create_and_trade_filter.normalize_dex_event(event),
906            DexEvent::PumpFunTrade(_)
907        ));
908    }
909
910    #[test]
911    fn all_protocol_groups_are_filterable() {
912        assert!(EventTypeFilter::include_only(vec![EventType::PumpFunTrade]).includes_pumpfun());
913        assert!(!EventTypeFilter::include_only(vec![EventType::AccountPumpFunGlobal])
914            .includes_pumpfun());
915        assert!(
916            !EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateAdmin]).includes_pumpfun()
917        );
918        assert!(EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]).includes_pumpswap());
919        assert!(EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateFeeShares])
920            .includes_pump_fees());
921        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumLaunchlabTrade])
922            .includes_raydium_launchlab());
923        assert!(
924            EventTypeFilter::include_only(vec![EventType::RaydiumCpmmSwap]).includes_raydium_cpmm()
925        );
926        assert!(
927            EventTypeFilter::include_only(vec![EventType::RaydiumClmmSwap]).includes_raydium_clmm()
928        );
929        assert!(!EventTypeFilter::include_only(vec![EventType::AccountRaydiumClmmPoolState])
930            .includes_raydium_clmm());
931        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumAmmV4Swap])
932            .includes_raydium_amm_v4());
933        assert!(EventTypeFilter::include_only(vec![EventType::OrcaWhirlpoolSwap])
934            .includes_orca_whirlpool());
935        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraPoolsSwap])
936            .includes_meteora_pools());
937        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2Swap])
938            .includes_meteora_damm_v2());
939        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2InitializePool])
940            .includes_meteora_damm_v2());
941        assert!(
942            EventTypeFilter::include_only(vec![EventType::MeteoraDlmmSwap]).includes_meteora_dlmm()
943        );
944        assert!(
945            EventTypeFilter::include_only(vec![EventType::MeteoraDbcSwap]).includes_meteora_dbc()
946        );
947    }
948
949    #[test]
950    fn exclude_filters_do_not_skip_whole_protocol_groups() {
951        let raydium = EventTypeFilter::exclude_types(vec![EventType::RaydiumCpmmSwap]);
952        assert!(raydium.includes_raydium_cpmm());
953        assert!(!raydium.should_include(EventType::RaydiumCpmmSwap));
954        assert!(raydium.should_include(EventType::RaydiumCpmmDeposit));
955
956        let all_cpmm = EventTypeFilter::exclude_types(vec![
957            EventType::RaydiumCpmmSwap,
958            EventType::RaydiumCpmmDeposit,
959            EventType::RaydiumCpmmWithdraw,
960            EventType::RaydiumCpmmInitialize,
961        ]);
962        assert!(!all_cpmm.includes_raydium_cpmm());
963
964        let all_launchlab = EventTypeFilter::exclude_types(vec![
965            EventType::RaydiumLaunchlabTrade,
966            EventType::RaydiumLaunchlabPoolCreate,
967            EventType::RaydiumLaunchlabMigrateAmm,
968        ]);
969        assert!(!all_launchlab.includes_raydium_launchlab());
970
971        let pump = EventTypeFilter::exclude_types(vec![EventType::PumpFunBuy]);
972        assert!(pump.includes_pumpfun());
973        assert!(!pump.should_include(EventType::PumpFunBuy));
974        assert!(!pump.should_include(EventType::PumpFunBuyExactSolIn));
975        assert!(pump.should_include(EventType::PumpFunSell));
976    }
977}
978
979#[derive(Debug, Clone)]
980pub struct SlotFilter {
981    pub min_slot: Option<u64>,
982    pub max_slot: Option<u64>,
983}
984
985impl SlotFilter {
986    pub fn new() -> Self {
987        Self { min_slot: None, max_slot: None }
988    }
989
990    pub fn min_slot(mut self, slot: u64) -> Self {
991        self.min_slot = Some(slot);
992        self
993    }
994
995    pub fn max_slot(mut self, slot: u64) -> Self {
996        self.max_slot = Some(slot);
997        self
998    }
999}
1000
1001impl Default for SlotFilter {
1002    fn default() -> Self {
1003        Self::new()
1004    }
1005}