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    /// All traffic from the shared LaunchLab program.
228    LaunchLab,
229    /// LaunchLab traffic attributed to StonkFun after parsing `platform_config`.
230    StonkFun,
231    /// Backward-compatible alias for [`Protocol::LaunchLab`].
232    RaydiumLaunchlab,
233    RaydiumCpmm,
234    RaydiumClmm,
235    RaydiumAmmV4,
236    OrcaWhirlpool,
237    MeteoraPools,
238    MeteoraDammV2,
239    MeteoraDlmm,
240    MeteoraDbc,
241}
242
243#[derive(Debug, Clone, Copy, PartialEq, Eq)]
244#[non_exhaustive]
245pub enum EventType {
246    // Block events
247    BlockMeta,
248
249    // RaydiumLaunchlab events
250    RaydiumLaunchlabTrade,
251    RaydiumLaunchlabPoolCreate,
252    RaydiumLaunchlabMigrateAmm,
253
254    // PumpFun events
255    PumpFunTrade,         // All trade events (backward compatible)
256    PumpFunBuy,           // Buy events only (filter by ix_name)
257    PumpFunSell,          // Sell events only (filter by ix_name)
258    PumpFunBuyExactSolIn, // BuyExactSolIn events only (filter by ix_name)
259    PumpFunCreate,
260    PumpFunCreateV2, // SPL-22 / Mayhem create
261    PumpFunComplete,
262    PumpFunMigrate,
263    /// Pump fees(`pfeeUx...`,`idls/pump_fees.json` Program data events)
264    PumpFeesCreateFeeSharingConfig,
265    PumpFeesInitializeFeeConfig,
266    PumpFeesResetFeeSharingConfig,
267    PumpFeesRevokeFeeSharingAuthority,
268    PumpFeesTransferFeeSharingAuthority,
269    PumpFeesUpdateAdmin,
270    PumpFeesUpdateFeeConfig,
271    PumpFeesUpdateFeeShares,
272    PumpFeesUpsertFeeTiers,
273    /// Pump.fun:`migrateBondingCurveCreatorEvent`
274    PumpFunMigrateBondingCurveCreator,
275
276    // PumpSwap events
277    PumpSwapTrade,
278    PumpSwapBuy,
279    PumpSwapSell,
280    PumpSwapCreatePool,
281    PumpSwapLiquidityAdded,
282    PumpSwapLiquidityRemoved,
283    // PumpSwapPoolUpdated,
284    // PumpSwapFeesClaimed,
285
286    // Raydium CPMM events
287    RaydiumCpmmSwap,
288    RaydiumCpmmDeposit,
289    RaydiumCpmmWithdraw,
290    RaydiumCpmmInitialize,
291
292    // Raydium CLMM events
293    RaydiumClmmSwap,
294    RaydiumClmmCreatePool,
295    RaydiumClmmOpenPosition,
296    RaydiumClmmClosePosition,
297    RaydiumClmmIncreaseLiquidity,
298    RaydiumClmmDecreaseLiquidity,
299    RaydiumClmmLiquidityChange,
300    RaydiumClmmConfigChange,
301    RaydiumClmmCreatePersonalPosition,
302    RaydiumClmmLiquidityCalculate,
303    RaydiumClmmOpenLimitOrder,
304    RaydiumClmmIncreaseLimitOrder,
305    RaydiumClmmDecreaseLimitOrder,
306    RaydiumClmmSettleLimitOrder,
307    RaydiumClmmUpdateRewardInfos,
308    RaydiumClmmOpenPositionWithTokenExtNft,
309    RaydiumClmmCollectFee,
310
311    // Raydium AMM V4 events
312    RaydiumAmmV4Swap,
313    RaydiumAmmV4Deposit,
314    RaydiumAmmV4Withdraw,
315    RaydiumAmmV4Initialize2,
316    RaydiumAmmV4WithdrawPnl,
317
318    // Orca Whirlpool events
319    OrcaWhirlpoolSwap,
320    OrcaWhirlpoolLiquidityIncreased,
321    OrcaWhirlpoolLiquidityDecreased,
322    OrcaWhirlpoolPoolInitialized,
323
324    // Meteora events
325    MeteoraPoolsSwap,
326    MeteoraPoolsAddLiquidity,
327    MeteoraPoolsRemoveLiquidity,
328    MeteoraPoolsBootstrapLiquidity,
329    MeteoraPoolsPoolCreated,
330    MeteoraPoolsSetPoolFees,
331
332    // Meteora DAMM V2 events
333    MeteoraDammV2Swap,
334    MeteoraDammV2AddLiquidity,
335    MeteoraDammV2RemoveLiquidity,
336    MeteoraDammV2InitializePool,
337    MeteoraDammV2CreatePosition,
338    MeteoraDammV2ClosePosition,
339    MeteoraDammV2UpdateDelegatePermission,
340    MeteoraDammV2WithdrawDeadLiquidityReward,
341    MeteoraDammV2CreateConfig,
342    MeteoraDammV2CreateDynamicConfig,
343    // MeteoraDammV2ClaimPositionFee,
344    // MeteoraDammV2InitializeReward,
345    // MeteoraDammV2FundReward,
346    // MeteoraDammV2ClaimReward,
347
348    // Meteora DBC events
349    MeteoraDbcSwap,
350    MeteoraDbcInitializePool,
351    MeteoraDbcCurveComplete,
352
353    // Meteora DLMM events
354    MeteoraDlmmSwap,
355    MeteoraDlmmAddLiquidity,
356    MeteoraDlmmRemoveLiquidity,
357    MeteoraDlmmInitializePool,
358    MeteoraDlmmInitializeBinArray,
359    MeteoraDlmmCreatePosition,
360    MeteoraDlmmClosePosition,
361    MeteoraDlmmClaimFee,
362
363    // Account events
364    TokenAccount,
365    TokenInfo,
366    NonceAccount,
367    AccountPumpFunGlobal,
368    AccountPumpFunBondingCurve,
369    AccountPumpFunFeeConfig,
370    AccountPumpFunSharingConfig,
371    AccountPumpFunGlobalVolumeAccumulator,
372    AccountPumpFunUserVolumeAccumulator,
373
374    AccountPumpSwapGlobalConfig,
375    AccountPumpSwapPool,
376    AccountRaydiumClmmAmmConfig,
377    AccountRaydiumClmmPoolState,
378    AccountRaydiumClmmTickArrayState,
379    AccountRaydiumCpmmAmmConfig,
380    AccountRaydiumCpmmPoolState,
381    AccountOrcaWhirlpool,
382    AccountOrcaPosition,
383    AccountOrcaTickArray,
384    AccountOrcaFeeTier,
385    AccountOrcaWhirlpoolsConfig,
386}
387
388#[derive(Debug, Clone)]
389pub struct EventTypeFilter {
390    pub include_only: Option<Vec<EventType>>,
391    pub exclude_types: Option<Vec<EventType>>,
392}
393
394impl EventTypeFilter {
395    pub fn include_only(types: Vec<EventType>) -> Self {
396        Self { include_only: Some(types), exclude_types: None }
397    }
398
399    pub fn exclude_types(types: Vec<EventType>) -> Self {
400        Self { include_only: None, exclude_types: Some(types) }
401    }
402
403    #[inline]
404    fn includes_any(&self, event_types: &[EventType]) -> bool {
405        event_types.iter().any(|event_type| self.should_include(*event_type))
406    }
407
408    pub fn should_include(&self, event_type: EventType) -> bool {
409        if let Some(ref include_only) = self.include_only {
410            // Direct match
411            if include_only.contains(&event_type) {
412                return true;
413            }
414            if matches!(
415                event_type,
416                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
417            ) {
418                if pumpfun_trade_filter_is_generic(include_only) {
419                    return true;
420                }
421                if event_type == EventType::PumpFunBuyExactSolIn
422                    && pumpfun_buy_filter_is_generic(include_only)
423                {
424                    return true;
425                }
426                return false;
427            }
428            if is_pumpfun_create_family(event_type) {
429                return include_only.iter().any(|t| is_pumpfun_create_family(*t));
430            }
431            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell) {
432                return include_only.contains(&EventType::PumpSwapTrade);
433            }
434            return false;
435        }
436
437        if let Some(ref exclude_types) = self.exclude_types {
438            if exclude_types.contains(&event_type) {
439                return false;
440            }
441            if matches!(
442                event_type,
443                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
444            ) && exclude_types.contains(&EventType::PumpFunTrade)
445            {
446                return false;
447            }
448            if event_type == EventType::PumpFunBuyExactSolIn
449                && exclude_types.contains(&EventType::PumpFunBuy)
450            {
451                return false;
452            }
453            if is_pumpfun_create_family(event_type)
454                && exclude_types.iter().any(|t| is_pumpfun_create_family(*t))
455            {
456                return false;
457            }
458            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell)
459                && exclude_types.contains(&EventType::PumpSwapTrade)
460            {
461                return false;
462            }
463            return true;
464        }
465
466        true
467    }
468
469    pub fn should_include_dex_event(&self, event: &crate::core::events::DexEvent) -> bool {
470        let Some(event_type) = event_type_from_dex_event(event) else { return true };
471        self.should_include(event_type)
472    }
473
474    #[inline]
475    pub fn includes_block_meta(&self) -> bool {
476        if let Some(ref include_only) = self.include_only {
477            return include_only.contains(&EventType::BlockMeta);
478        }
479        false
480    }
481
482    #[inline]
483    pub fn normalize_dex_event(
484        &self,
485        event: crate::core::events::DexEvent,
486    ) -> crate::core::events::DexEvent {
487        use crate::core::events::DexEvent;
488
489        let Some(ref include_only) = self.include_only else { return event };
490        if pumpfun_trade_filter_is_generic(include_only) {
491            return match event {
492                DexEvent::PumpFunBuy(t)
493                | DexEvent::PumpFunSell(t)
494                | DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunTrade(t),
495                other => other,
496            };
497        }
498        if pumpfun_buy_filter_is_generic(include_only) {
499            return match event {
500                DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunBuy(t),
501                other => other,
502            };
503        }
504
505        event
506    }
507
508    #[inline]
509    pub fn includes_pumpfun(&self) -> bool {
510        self.includes_any(&[
511            EventType::PumpFunTrade,
512            EventType::PumpFunBuy,
513            EventType::PumpFunSell,
514            EventType::PumpFunBuyExactSolIn,
515            EventType::PumpFunCreate,
516            EventType::PumpFunCreateV2,
517            EventType::PumpFunComplete,
518            EventType::PumpFunMigrate,
519            EventType::PumpFunMigrateBondingCurveCreator,
520        ])
521    }
522
523    #[inline]
524    pub fn includes_meteora_damm_v2(&self) -> bool {
525        self.includes_any(&[
526            EventType::MeteoraDammV2Swap,
527            EventType::MeteoraDammV2AddLiquidity,
528            EventType::MeteoraDammV2CreatePosition,
529            EventType::MeteoraDammV2ClosePosition,
530            EventType::MeteoraDammV2InitializePool,
531            EventType::MeteoraDammV2RemoveLiquidity,
532            EventType::MeteoraDammV2UpdateDelegatePermission,
533            EventType::MeteoraDammV2WithdrawDeadLiquidityReward,
534            EventType::MeteoraDammV2CreateConfig,
535            EventType::MeteoraDammV2CreateDynamicConfig,
536        ])
537    }
538
539    #[inline]
540    pub fn includes_pump_fees(&self) -> bool {
541        self.includes_any(&[
542            EventType::PumpFeesCreateFeeSharingConfig,
543            EventType::PumpFeesInitializeFeeConfig,
544            EventType::PumpFeesResetFeeSharingConfig,
545            EventType::PumpFeesRevokeFeeSharingAuthority,
546            EventType::PumpFeesTransferFeeSharingAuthority,
547            EventType::PumpFeesUpdateAdmin,
548            EventType::PumpFeesUpdateFeeConfig,
549            EventType::PumpFeesUpdateFeeShares,
550            EventType::PumpFeesUpsertFeeTiers,
551        ])
552    }
553
554    /// Check if PumpSwap protocol events are included in the filter
555    #[inline]
556    pub fn includes_pumpswap(&self) -> bool {
557        self.includes_any(&[
558            EventType::PumpSwapTrade,
559            EventType::PumpSwapBuy,
560            EventType::PumpSwapSell,
561            EventType::PumpSwapCreatePool,
562            EventType::PumpSwapLiquidityAdded,
563            EventType::PumpSwapLiquidityRemoved,
564        ])
565    }
566
567    /// Check if LaunchLab events are included in the filter.
568    #[inline]
569    pub fn includes_raydium_launchlab(&self) -> bool {
570        self.includes_any(&[
571            EventType::RaydiumLaunchlabTrade,
572            EventType::RaydiumLaunchlabPoolCreate,
573            EventType::RaydiumLaunchlabMigrateAmm,
574        ])
575    }
576
577    #[inline]
578    pub fn includes_raydium_cpmm(&self) -> bool {
579        self.includes_any(&[
580            EventType::RaydiumCpmmSwap,
581            EventType::RaydiumCpmmDeposit,
582            EventType::RaydiumCpmmWithdraw,
583            EventType::RaydiumCpmmInitialize,
584        ])
585    }
586
587    #[inline]
588    pub fn includes_raydium_clmm(&self) -> bool {
589        self.includes_any(&[
590            EventType::RaydiumClmmSwap,
591            EventType::RaydiumClmmCreatePool,
592            EventType::RaydiumClmmOpenPosition,
593            EventType::RaydiumClmmClosePosition,
594            EventType::RaydiumClmmIncreaseLiquidity,
595            EventType::RaydiumClmmDecreaseLiquidity,
596            EventType::RaydiumClmmLiquidityChange,
597            EventType::RaydiumClmmConfigChange,
598            EventType::RaydiumClmmCreatePersonalPosition,
599            EventType::RaydiumClmmLiquidityCalculate,
600            EventType::RaydiumClmmOpenLimitOrder,
601            EventType::RaydiumClmmIncreaseLimitOrder,
602            EventType::RaydiumClmmDecreaseLimitOrder,
603            EventType::RaydiumClmmSettleLimitOrder,
604            EventType::RaydiumClmmUpdateRewardInfos,
605            EventType::RaydiumClmmOpenPositionWithTokenExtNft,
606            EventType::RaydiumClmmCollectFee,
607        ])
608    }
609
610    #[inline]
611    pub fn includes_raydium_amm_v4(&self) -> bool {
612        self.includes_any(&[
613            EventType::RaydiumAmmV4Swap,
614            EventType::RaydiumAmmV4Deposit,
615            EventType::RaydiumAmmV4Withdraw,
616            EventType::RaydiumAmmV4Initialize2,
617            EventType::RaydiumAmmV4WithdrawPnl,
618        ])
619    }
620
621    #[inline]
622    pub fn includes_orca_whirlpool(&self) -> bool {
623        self.includes_any(&[
624            EventType::OrcaWhirlpoolSwap,
625            EventType::OrcaWhirlpoolLiquidityIncreased,
626            EventType::OrcaWhirlpoolLiquidityDecreased,
627            EventType::OrcaWhirlpoolPoolInitialized,
628        ])
629    }
630
631    #[inline]
632    pub fn includes_meteora_pools(&self) -> bool {
633        self.includes_any(&[
634            EventType::MeteoraPoolsSwap,
635            EventType::MeteoraPoolsAddLiquidity,
636            EventType::MeteoraPoolsRemoveLiquidity,
637            EventType::MeteoraPoolsBootstrapLiquidity,
638            EventType::MeteoraPoolsPoolCreated,
639            EventType::MeteoraPoolsSetPoolFees,
640        ])
641    }
642
643    #[inline]
644    pub fn includes_meteora_dlmm(&self) -> bool {
645        self.includes_any(&[
646            EventType::MeteoraDlmmSwap,
647            EventType::MeteoraDlmmAddLiquidity,
648            EventType::MeteoraDlmmRemoveLiquidity,
649            EventType::MeteoraDlmmInitializePool,
650            EventType::MeteoraDlmmInitializeBinArray,
651            EventType::MeteoraDlmmCreatePosition,
652            EventType::MeteoraDlmmClosePosition,
653            EventType::MeteoraDlmmClaimFee,
654        ])
655    }
656
657    #[inline]
658    pub fn includes_meteora_dbc(&self) -> bool {
659        self.includes_any(&[
660            EventType::MeteoraDbcSwap,
661            EventType::MeteoraDbcInitializePool,
662            EventType::MeteoraDbcCurveComplete,
663        ])
664    }
665}
666
667#[inline]
668fn pumpfun_trade_filter_is_generic(include_only: &[EventType]) -> bool {
669    include_only.contains(&EventType::PumpFunTrade)
670        && !include_only.iter().any(|t| {
671            matches!(
672                t,
673                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
674            )
675        })
676}
677
678#[inline]
679fn pumpfun_buy_filter_is_generic(include_only: &[EventType]) -> bool {
680    include_only.contains(&EventType::PumpFunBuy)
681        && !include_only.contains(&EventType::PumpFunBuyExactSolIn)
682}
683
684#[inline]
685fn is_pumpfun_create_family(event_type: EventType) -> bool {
686    matches!(event_type, EventType::PumpFunCreate | EventType::PumpFunCreateV2)
687}
688
689#[inline]
690pub fn event_type_from_dex_event(event: &crate::core::events::DexEvent) -> Option<EventType> {
691    use crate::core::events::DexEvent;
692    match event {
693        DexEvent::PumpFunCreate(_) => Some(EventType::PumpFunCreate),
694        DexEvent::PumpFunCreateV2(_) => Some(EventType::PumpFunCreateV2),
695        DexEvent::PumpFunTrade(_) => Some(EventType::PumpFunTrade),
696        DexEvent::PumpFunBuy(_) => Some(EventType::PumpFunBuy),
697        DexEvent::PumpFunSell(_) => Some(EventType::PumpFunSell),
698        DexEvent::PumpFunBuyExactSolIn(_) => Some(EventType::PumpFunBuyExactSolIn),
699        DexEvent::PumpFunMigrate(_) => Some(EventType::PumpFunMigrate),
700        DexEvent::PumpFeesCreateFeeSharingConfig(_) => {
701            Some(EventType::PumpFeesCreateFeeSharingConfig)
702        }
703        DexEvent::PumpFeesInitializeFeeConfig(_) => Some(EventType::PumpFeesInitializeFeeConfig),
704        DexEvent::PumpFeesResetFeeSharingConfig(_) => {
705            Some(EventType::PumpFeesResetFeeSharingConfig)
706        }
707        DexEvent::PumpFeesRevokeFeeSharingAuthority(_) => {
708            Some(EventType::PumpFeesRevokeFeeSharingAuthority)
709        }
710        DexEvent::PumpFeesTransferFeeSharingAuthority(_) => {
711            Some(EventType::PumpFeesTransferFeeSharingAuthority)
712        }
713        DexEvent::PumpFeesUpdateAdmin(_) => Some(EventType::PumpFeesUpdateAdmin),
714        DexEvent::PumpFeesUpdateFeeConfig(_) => Some(EventType::PumpFeesUpdateFeeConfig),
715        DexEvent::PumpFeesUpdateFeeShares(_) => Some(EventType::PumpFeesUpdateFeeShares),
716        DexEvent::PumpFeesUpsertFeeTiers(_) => Some(EventType::PumpFeesUpsertFeeTiers),
717        DexEvent::PumpFunMigrateBondingCurveCreator(_) => {
718            Some(EventType::PumpFunMigrateBondingCurveCreator)
719        }
720        DexEvent::PumpFunGlobalAccount(_) => Some(EventType::AccountPumpFunGlobal),
721        DexEvent::PumpFunBondingCurveAccount(_) => Some(EventType::AccountPumpFunBondingCurve),
722        DexEvent::PumpFunFeeConfigAccount(_) => Some(EventType::AccountPumpFunFeeConfig),
723        DexEvent::PumpFunSharingConfigAccount(_) => Some(EventType::AccountPumpFunSharingConfig),
724        DexEvent::PumpFunGlobalVolumeAccumulatorAccount(_) => {
725            Some(EventType::AccountPumpFunGlobalVolumeAccumulator)
726        }
727        DexEvent::PumpFunUserVolumeAccumulatorAccount(_) => {
728            Some(EventType::AccountPumpFunUserVolumeAccumulator)
729        }
730        DexEvent::PumpSwapTrade(_) => Some(EventType::PumpSwapTrade),
731        DexEvent::PumpSwapBuy(_) => Some(EventType::PumpSwapBuy),
732        DexEvent::PumpSwapSell(_) => Some(EventType::PumpSwapSell),
733        DexEvent::PumpSwapCreatePool(_) => Some(EventType::PumpSwapCreatePool),
734        DexEvent::PumpSwapLiquidityAdded(_) => Some(EventType::PumpSwapLiquidityAdded),
735        DexEvent::PumpSwapLiquidityRemoved(_) => Some(EventType::PumpSwapLiquidityRemoved),
736        DexEvent::MeteoraDammV2Swap(_) => Some(EventType::MeteoraDammV2Swap),
737        DexEvent::MeteoraDammV2CreatePosition(_) => Some(EventType::MeteoraDammV2CreatePosition),
738        DexEvent::MeteoraDammV2ClosePosition(_) => Some(EventType::MeteoraDammV2ClosePosition),
739        DexEvent::MeteoraDammV2AddLiquidity(_) => Some(EventType::MeteoraDammV2AddLiquidity),
740        DexEvent::MeteoraDammV2RemoveLiquidity(_) => Some(EventType::MeteoraDammV2RemoveLiquidity),
741        DexEvent::MeteoraDammV2InitializePool(_) => Some(EventType::MeteoraDammV2InitializePool),
742        DexEvent::MeteoraDammV2UpdateDelegatePermission(_) => {
743            Some(EventType::MeteoraDammV2UpdateDelegatePermission)
744        }
745        DexEvent::MeteoraDammV2WithdrawDeadLiquidityReward(_) => {
746            Some(EventType::MeteoraDammV2WithdrawDeadLiquidityReward)
747        }
748        DexEvent::MeteoraDammV2CreateConfig(_) => Some(EventType::MeteoraDammV2CreateConfig),
749        DexEvent::MeteoraDammV2CreateDynamicConfig(_) => {
750            Some(EventType::MeteoraDammV2CreateDynamicConfig)
751        }
752        DexEvent::MeteoraDbcSwap(_) => Some(EventType::MeteoraDbcSwap),
753        DexEvent::MeteoraDbcInitializePool(_) => Some(EventType::MeteoraDbcInitializePool),
754        DexEvent::MeteoraDbcCurveComplete(_) => Some(EventType::MeteoraDbcCurveComplete),
755        DexEvent::RaydiumLaunchlabTrade(_) => Some(EventType::RaydiumLaunchlabTrade),
756        DexEvent::RaydiumLaunchlabPoolCreate(_) => Some(EventType::RaydiumLaunchlabPoolCreate),
757        DexEvent::RaydiumLaunchlabMigrateAmm(_) => Some(EventType::RaydiumLaunchlabMigrateAmm),
758        DexEvent::RaydiumClmmSwap(_) => Some(EventType::RaydiumClmmSwap),
759        DexEvent::RaydiumClmmCreatePool(_) => Some(EventType::RaydiumClmmCreatePool),
760        DexEvent::RaydiumClmmOpenPosition(_) => Some(EventType::RaydiumClmmOpenPosition),
761        DexEvent::RaydiumClmmOpenPositionWithTokenExtNft(_) => {
762            Some(EventType::RaydiumClmmOpenPositionWithTokenExtNft)
763        }
764        DexEvent::RaydiumClmmClosePosition(_) => Some(EventType::RaydiumClmmClosePosition),
765        DexEvent::RaydiumClmmIncreaseLiquidity(_) => Some(EventType::RaydiumClmmIncreaseLiquidity),
766        DexEvent::RaydiumClmmDecreaseLiquidity(_) => Some(EventType::RaydiumClmmDecreaseLiquidity),
767        DexEvent::RaydiumClmmLiquidityChange(_) => Some(EventType::RaydiumClmmLiquidityChange),
768        DexEvent::RaydiumClmmConfigChange(_) => Some(EventType::RaydiumClmmConfigChange),
769        DexEvent::RaydiumClmmCreatePersonalPosition(_) => {
770            Some(EventType::RaydiumClmmCreatePersonalPosition)
771        }
772        DexEvent::RaydiumClmmLiquidityCalculate(_) => {
773            Some(EventType::RaydiumClmmLiquidityCalculate)
774        }
775        DexEvent::RaydiumClmmOpenLimitOrder(_) => Some(EventType::RaydiumClmmOpenLimitOrder),
776        DexEvent::RaydiumClmmIncreaseLimitOrder(_) => {
777            Some(EventType::RaydiumClmmIncreaseLimitOrder)
778        }
779        DexEvent::RaydiumClmmDecreaseLimitOrder(_) => {
780            Some(EventType::RaydiumClmmDecreaseLimitOrder)
781        }
782        DexEvent::RaydiumClmmSettleLimitOrder(_) => Some(EventType::RaydiumClmmSettleLimitOrder),
783        DexEvent::RaydiumClmmUpdateRewardInfos(_) => Some(EventType::RaydiumClmmUpdateRewardInfos),
784        DexEvent::RaydiumClmmCollectFee(_) => Some(EventType::RaydiumClmmCollectFee),
785        DexEvent::RaydiumClmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumClmmAmmConfig),
786        DexEvent::RaydiumClmmPoolStateAccount(_) => Some(EventType::AccountRaydiumClmmPoolState),
787        DexEvent::RaydiumClmmTickArrayStateAccount(_) => {
788            Some(EventType::AccountRaydiumClmmTickArrayState)
789        }
790        DexEvent::RaydiumCpmmSwap(_) => Some(EventType::RaydiumCpmmSwap),
791        DexEvent::RaydiumCpmmDeposit(_) => Some(EventType::RaydiumCpmmDeposit),
792        DexEvent::RaydiumCpmmWithdraw(_) => Some(EventType::RaydiumCpmmWithdraw),
793        DexEvent::RaydiumCpmmInitialize(_) => Some(EventType::RaydiumCpmmInitialize),
794        DexEvent::RaydiumCpmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumCpmmAmmConfig),
795        DexEvent::RaydiumCpmmPoolStateAccount(_) => Some(EventType::AccountRaydiumCpmmPoolState),
796        DexEvent::RaydiumAmmV4Swap(_) => Some(EventType::RaydiumAmmV4Swap),
797        DexEvent::RaydiumAmmV4Deposit(_) => Some(EventType::RaydiumAmmV4Deposit),
798        DexEvent::RaydiumAmmV4Initialize2(_) => Some(EventType::RaydiumAmmV4Initialize2),
799        DexEvent::RaydiumAmmV4Withdraw(_) => Some(EventType::RaydiumAmmV4Withdraw),
800        DexEvent::RaydiumAmmV4WithdrawPnl(_) => Some(EventType::RaydiumAmmV4WithdrawPnl),
801        DexEvent::OrcaWhirlpoolSwap(_) => Some(EventType::OrcaWhirlpoolSwap),
802        DexEvent::OrcaWhirlpoolLiquidityIncreased(_) => {
803            Some(EventType::OrcaWhirlpoolLiquidityIncreased)
804        }
805        DexEvent::OrcaWhirlpoolLiquidityDecreased(_) => {
806            Some(EventType::OrcaWhirlpoolLiquidityDecreased)
807        }
808        DexEvent::OrcaWhirlpoolPoolInitialized(_) => Some(EventType::OrcaWhirlpoolPoolInitialized),
809        DexEvent::OrcaWhirlpoolAccount(_) => Some(EventType::AccountOrcaWhirlpool),
810        DexEvent::OrcaPositionAccount(_) => Some(EventType::AccountOrcaPosition),
811        DexEvent::OrcaTickArrayAccount(_) => Some(EventType::AccountOrcaTickArray),
812        DexEvent::OrcaFeeTierAccount(_) => Some(EventType::AccountOrcaFeeTier),
813        DexEvent::OrcaWhirlpoolsConfigAccount(_) => Some(EventType::AccountOrcaWhirlpoolsConfig),
814        DexEvent::MeteoraPoolsSwap(_) => Some(EventType::MeteoraPoolsSwap),
815        DexEvent::MeteoraPoolsAddLiquidity(_) => Some(EventType::MeteoraPoolsAddLiquidity),
816        DexEvent::MeteoraPoolsRemoveLiquidity(_) => Some(EventType::MeteoraPoolsRemoveLiquidity),
817        DexEvent::MeteoraPoolsBootstrapLiquidity(_) => {
818            Some(EventType::MeteoraPoolsBootstrapLiquidity)
819        }
820        DexEvent::MeteoraPoolsPoolCreated(_) => Some(EventType::MeteoraPoolsPoolCreated),
821        DexEvent::MeteoraPoolsSetPoolFees(_) => Some(EventType::MeteoraPoolsSetPoolFees),
822        DexEvent::MeteoraDlmmSwap(_) => Some(EventType::MeteoraDlmmSwap),
823        DexEvent::MeteoraDlmmAddLiquidity(_) => Some(EventType::MeteoraDlmmAddLiquidity),
824        DexEvent::MeteoraDlmmRemoveLiquidity(_) => Some(EventType::MeteoraDlmmRemoveLiquidity),
825        DexEvent::MeteoraDlmmInitializePool(_) => Some(EventType::MeteoraDlmmInitializePool),
826        DexEvent::MeteoraDlmmInitializeBinArray(_) => {
827            Some(EventType::MeteoraDlmmInitializeBinArray)
828        }
829        DexEvent::MeteoraDlmmCreatePosition(_) => Some(EventType::MeteoraDlmmCreatePosition),
830        DexEvent::MeteoraDlmmClosePosition(_) => Some(EventType::MeteoraDlmmClosePosition),
831        DexEvent::MeteoraDlmmClaimFee(_) => Some(EventType::MeteoraDlmmClaimFee),
832        DexEvent::TokenAccount(_) => Some(EventType::TokenAccount),
833        DexEvent::TokenInfo(_) => Some(EventType::TokenInfo),
834        DexEvent::NonceAccount(_) => Some(EventType::NonceAccount),
835        DexEvent::PumpSwapGlobalConfigAccount(_) => Some(EventType::AccountPumpSwapGlobalConfig),
836        DexEvent::PumpSwapPoolAccount(_) => Some(EventType::AccountPumpSwapPool),
837        DexEvent::BlockMeta(_) => Some(EventType::BlockMeta),
838        DexEvent::Error(_) => None,
839    }
840}
841
842#[cfg(test)]
843mod event_type_filter_tests {
844    use super::*;
845
846    #[test]
847    fn generic_trade_filters_cover_specific_trade_variants() {
848        let pump = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
849        assert!(pump.should_include(EventType::PumpFunTrade));
850        assert!(pump.should_include(EventType::PumpFunBuy));
851        assert!(pump.should_include(EventType::PumpFunSell));
852        assert!(pump.should_include(EventType::PumpFunBuyExactSolIn));
853
854        let pump_specific = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
855        assert!(!pump_specific.should_include(EventType::PumpFunTrade));
856        assert!(pump_specific.should_include(EventType::PumpFunBuy));
857        assert!(!pump_specific.should_include(EventType::PumpFunSell));
858        assert!(pump_specific.should_include(EventType::PumpFunBuyExactSolIn));
859
860        let pump_exact_buy = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
861        assert!(!pump_exact_buy.should_include(EventType::PumpFunTrade));
862        assert!(!pump_exact_buy.should_include(EventType::PumpFunBuy));
863        assert!(!pump_exact_buy.should_include(EventType::PumpFunSell));
864        assert!(pump_exact_buy.should_include(EventType::PumpFunBuyExactSolIn));
865
866        let pumpswap = EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]);
867        assert!(pumpswap.should_include(EventType::PumpSwapBuy));
868        assert!(pumpswap.should_include(EventType::PumpSwapSell));
869
870        let exclude_pumpswap = EventTypeFilter::exclude_types(vec![EventType::PumpSwapTrade]);
871        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapBuy));
872        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapSell));
873    }
874
875    #[test]
876    fn generic_pumpfun_trade_filter_normalizes_specific_variants() {
877        use crate::core::events::{DexEvent, PumpFunTradeEvent};
878
879        let filter = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
880        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
881        assert!(matches!(filter.normalize_dex_event(event), DexEvent::PumpFunTrade(_)));
882
883        let specific_filter =
884            EventTypeFilter::include_only(vec![EventType::PumpFunTrade, EventType::PumpFunBuy]);
885        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
886        assert!(matches!(specific_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
887
888        let buy_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
889        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
890            is_buy: true,
891            ..Default::default()
892        });
893        assert!(matches!(buy_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
894
895        let exact_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
896        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
897            is_buy: true,
898            ..Default::default()
899        });
900        assert!(matches!(
901            exact_filter.normalize_dex_event(event),
902            DexEvent::PumpFunBuyExactSolIn(_)
903        ));
904
905        let create_and_trade_filter =
906            EventTypeFilter::include_only(vec![EventType::PumpFunCreate, EventType::PumpFunTrade]);
907        let event =
908            DexEvent::PumpFunSell(PumpFunTradeEvent { is_buy: false, ..Default::default() });
909        assert!(matches!(
910            create_and_trade_filter.normalize_dex_event(event),
911            DexEvent::PumpFunTrade(_)
912        ));
913    }
914
915    #[test]
916    fn all_protocol_groups_are_filterable() {
917        assert!(EventTypeFilter::include_only(vec![EventType::PumpFunTrade]).includes_pumpfun());
918        assert!(!EventTypeFilter::include_only(vec![EventType::AccountPumpFunGlobal])
919            .includes_pumpfun());
920        assert!(
921            !EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateAdmin]).includes_pumpfun()
922        );
923        assert!(EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]).includes_pumpswap());
924        assert!(EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateFeeShares])
925            .includes_pump_fees());
926        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumLaunchlabTrade])
927            .includes_raydium_launchlab());
928        assert!(
929            EventTypeFilter::include_only(vec![EventType::RaydiumCpmmSwap]).includes_raydium_cpmm()
930        );
931        assert!(
932            EventTypeFilter::include_only(vec![EventType::RaydiumClmmSwap]).includes_raydium_clmm()
933        );
934        assert!(!EventTypeFilter::include_only(vec![EventType::AccountRaydiumClmmPoolState])
935            .includes_raydium_clmm());
936        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumAmmV4Swap])
937            .includes_raydium_amm_v4());
938        assert!(EventTypeFilter::include_only(vec![EventType::OrcaWhirlpoolSwap])
939            .includes_orca_whirlpool());
940        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraPoolsSwap])
941            .includes_meteora_pools());
942        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2Swap])
943            .includes_meteora_damm_v2());
944        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2InitializePool])
945            .includes_meteora_damm_v2());
946        assert!(
947            EventTypeFilter::include_only(vec![EventType::MeteoraDlmmSwap]).includes_meteora_dlmm()
948        );
949        assert!(
950            EventTypeFilter::include_only(vec![EventType::MeteoraDbcSwap]).includes_meteora_dbc()
951        );
952    }
953
954    #[test]
955    fn exclude_filters_do_not_skip_whole_protocol_groups() {
956        let raydium = EventTypeFilter::exclude_types(vec![EventType::RaydiumCpmmSwap]);
957        assert!(raydium.includes_raydium_cpmm());
958        assert!(!raydium.should_include(EventType::RaydiumCpmmSwap));
959        assert!(raydium.should_include(EventType::RaydiumCpmmDeposit));
960
961        let all_cpmm = EventTypeFilter::exclude_types(vec![
962            EventType::RaydiumCpmmSwap,
963            EventType::RaydiumCpmmDeposit,
964            EventType::RaydiumCpmmWithdraw,
965            EventType::RaydiumCpmmInitialize,
966        ]);
967        assert!(!all_cpmm.includes_raydium_cpmm());
968
969        let all_launchlab = EventTypeFilter::exclude_types(vec![
970            EventType::RaydiumLaunchlabTrade,
971            EventType::RaydiumLaunchlabPoolCreate,
972            EventType::RaydiumLaunchlabMigrateAmm,
973        ]);
974        assert!(!all_launchlab.includes_raydium_launchlab());
975
976        let pump = EventTypeFilter::exclude_types(vec![EventType::PumpFunBuy]);
977        assert!(pump.includes_pumpfun());
978        assert!(!pump.should_include(EventType::PumpFunBuy));
979        assert!(!pump.should_include(EventType::PumpFunBuyExactSolIn));
980        assert!(pump.should_include(EventType::PumpFunSell));
981    }
982}
983
984#[derive(Debug, Clone)]
985pub struct SlotFilter {
986    pub min_slot: Option<u64>,
987    pub max_slot: Option<u64>,
988}
989
990impl SlotFilter {
991    pub fn new() -> Self {
992        Self { min_slot: None, max_slot: None }
993    }
994
995    pub fn min_slot(mut self, slot: u64) -> Self {
996        self.min_slot = Some(slot);
997        self
998    }
999
1000    pub fn max_slot(mut self, slot: u64) -> Self {
1001        self.max_slot = Some(slot);
1002        self
1003    }
1004}
1005
1006impl Default for SlotFilter {
1007    fn default() -> Self {
1008        Self::new()
1009    }
1010}