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