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    AccountLiquiditySnapshot,
385    /// Raw subscription bytes with slot/write_version; explicitly opt in via include_only.
386    AccountRawSnapshot,
387    AccountOrcaFeeTier,
388    AccountOrcaWhirlpoolsConfig,
389}
390
391#[derive(Debug, Clone)]
392pub struct EventTypeFilter {
393    pub include_only: Option<Vec<EventType>>,
394    pub exclude_types: Option<Vec<EventType>>,
395}
396
397impl EventTypeFilter {
398    pub fn include_only(types: Vec<EventType>) -> Self {
399        Self { include_only: Some(types), exclude_types: None }
400    }
401
402    pub fn exclude_types(types: Vec<EventType>) -> Self {
403        Self { include_only: None, exclude_types: Some(types) }
404    }
405
406    #[inline]
407    fn includes_any(&self, event_types: &[EventType]) -> bool {
408        event_types.iter().any(|event_type| self.should_include(*event_type))
409    }
410
411    pub fn should_include(&self, event_type: EventType) -> bool {
412        if let Some(ref include_only) = self.include_only {
413            // Direct match
414            if include_only.contains(&event_type) {
415                return true;
416            }
417            if matches!(
418                event_type,
419                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
420            ) {
421                if pumpfun_trade_filter_is_generic(include_only) {
422                    return true;
423                }
424                if event_type == EventType::PumpFunBuyExactSolIn
425                    && pumpfun_buy_filter_is_generic(include_only)
426                {
427                    return true;
428                }
429                return false;
430            }
431            if is_pumpfun_create_family(event_type) {
432                return include_only.iter().any(|t| is_pumpfun_create_family(*t));
433            }
434            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell) {
435                return include_only.contains(&EventType::PumpSwapTrade);
436            }
437            return false;
438        }
439
440        if let Some(ref exclude_types) = self.exclude_types {
441            if exclude_types.contains(&event_type) {
442                return false;
443            }
444            if matches!(
445                event_type,
446                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
447            ) && exclude_types.contains(&EventType::PumpFunTrade)
448            {
449                return false;
450            }
451            if event_type == EventType::PumpFunBuyExactSolIn
452                && exclude_types.contains(&EventType::PumpFunBuy)
453            {
454                return false;
455            }
456            if is_pumpfun_create_family(event_type)
457                && exclude_types.iter().any(|t| is_pumpfun_create_family(*t))
458            {
459                return false;
460            }
461            if matches!(event_type, EventType::PumpSwapBuy | EventType::PumpSwapSell)
462                && exclude_types.contains(&EventType::PumpSwapTrade)
463            {
464                return false;
465            }
466            return true;
467        }
468
469        true
470    }
471
472    pub fn should_include_dex_event(&self, event: &crate::core::events::DexEvent) -> bool {
473        let Some(event_type) = event_type_from_dex_event(event) else { return true };
474        self.should_include(event_type)
475    }
476
477    #[inline]
478    pub fn includes_block_meta(&self) -> bool {
479        if let Some(ref include_only) = self.include_only {
480            return include_only.contains(&EventType::BlockMeta);
481        }
482        false
483    }
484
485    #[inline]
486    pub fn normalize_dex_event(
487        &self,
488        event: crate::core::events::DexEvent,
489    ) -> crate::core::events::DexEvent {
490        use crate::core::events::DexEvent;
491
492        let Some(ref include_only) = self.include_only else { return event };
493        if pumpfun_trade_filter_is_generic(include_only) {
494            return match event {
495                DexEvent::PumpFunBuy(t)
496                | DexEvent::PumpFunSell(t)
497                | DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunTrade(t),
498                other => other,
499            };
500        }
501        if pumpfun_buy_filter_is_generic(include_only) {
502            return match event {
503                DexEvent::PumpFunBuyExactSolIn(t) => DexEvent::PumpFunBuy(t),
504                other => other,
505            };
506        }
507
508        event
509    }
510
511    #[inline]
512    pub fn includes_pumpfun(&self) -> bool {
513        self.includes_any(&[
514            EventType::PumpFunTrade,
515            EventType::PumpFunBuy,
516            EventType::PumpFunSell,
517            EventType::PumpFunBuyExactSolIn,
518            EventType::PumpFunCreate,
519            EventType::PumpFunCreateV2,
520            EventType::PumpFunComplete,
521            EventType::PumpFunMigrate,
522            EventType::PumpFunMigrateBondingCurveCreator,
523        ])
524    }
525
526    #[inline]
527    pub fn includes_meteora_damm_v2(&self) -> bool {
528        self.includes_any(&[
529            EventType::MeteoraDammV2Swap,
530            EventType::MeteoraDammV2AddLiquidity,
531            EventType::MeteoraDammV2CreatePosition,
532            EventType::MeteoraDammV2ClosePosition,
533            EventType::MeteoraDammV2InitializePool,
534            EventType::MeteoraDammV2RemoveLiquidity,
535            EventType::MeteoraDammV2UpdateDelegatePermission,
536            EventType::MeteoraDammV2WithdrawDeadLiquidityReward,
537            EventType::MeteoraDammV2CreateConfig,
538            EventType::MeteoraDammV2CreateDynamicConfig,
539        ])
540    }
541
542    #[inline]
543    pub fn includes_pump_fees(&self) -> bool {
544        self.includes_any(&[
545            EventType::PumpFeesCreateFeeSharingConfig,
546            EventType::PumpFeesInitializeFeeConfig,
547            EventType::PumpFeesResetFeeSharingConfig,
548            EventType::PumpFeesRevokeFeeSharingAuthority,
549            EventType::PumpFeesTransferFeeSharingAuthority,
550            EventType::PumpFeesUpdateAdmin,
551            EventType::PumpFeesUpdateFeeConfig,
552            EventType::PumpFeesUpdateFeeShares,
553            EventType::PumpFeesUpsertFeeTiers,
554        ])
555    }
556
557    /// Check if PumpSwap protocol events are included in the filter
558    #[inline]
559    pub fn includes_pumpswap(&self) -> bool {
560        self.includes_any(&[
561            EventType::PumpSwapTrade,
562            EventType::PumpSwapBuy,
563            EventType::PumpSwapSell,
564            EventType::PumpSwapCreatePool,
565            EventType::PumpSwapLiquidityAdded,
566            EventType::PumpSwapLiquidityRemoved,
567        ])
568    }
569
570    /// Check if LaunchLab events are included in the filter.
571    #[inline]
572    pub fn includes_raydium_launchlab(&self) -> bool {
573        self.includes_any(&[
574            EventType::RaydiumLaunchlabTrade,
575            EventType::RaydiumLaunchlabPoolCreate,
576            EventType::RaydiumLaunchlabMigrateAmm,
577        ])
578    }
579
580    #[inline]
581    pub fn includes_raydium_cpmm(&self) -> bool {
582        self.includes_any(&[
583            EventType::RaydiumCpmmSwap,
584            EventType::RaydiumCpmmDeposit,
585            EventType::RaydiumCpmmWithdraw,
586            EventType::RaydiumCpmmInitialize,
587        ])
588    }
589
590    #[inline]
591    pub fn includes_raydium_clmm(&self) -> bool {
592        self.includes_any(&[
593            EventType::RaydiumClmmSwap,
594            EventType::RaydiumClmmCreatePool,
595            EventType::RaydiumClmmOpenPosition,
596            EventType::RaydiumClmmClosePosition,
597            EventType::RaydiumClmmIncreaseLiquidity,
598            EventType::RaydiumClmmDecreaseLiquidity,
599            EventType::RaydiumClmmLiquidityChange,
600            EventType::RaydiumClmmConfigChange,
601            EventType::RaydiumClmmCreatePersonalPosition,
602            EventType::RaydiumClmmLiquidityCalculate,
603            EventType::RaydiumClmmOpenLimitOrder,
604            EventType::RaydiumClmmIncreaseLimitOrder,
605            EventType::RaydiumClmmDecreaseLimitOrder,
606            EventType::RaydiumClmmSettleLimitOrder,
607            EventType::RaydiumClmmUpdateRewardInfos,
608            EventType::RaydiumClmmOpenPositionWithTokenExtNft,
609            EventType::RaydiumClmmCollectFee,
610        ])
611    }
612
613    #[inline]
614    pub fn includes_raydium_amm_v4(&self) -> bool {
615        self.includes_any(&[
616            EventType::RaydiumAmmV4Swap,
617            EventType::RaydiumAmmV4Deposit,
618            EventType::RaydiumAmmV4Withdraw,
619            EventType::RaydiumAmmV4Initialize2,
620            EventType::RaydiumAmmV4WithdrawPnl,
621        ])
622    }
623
624    #[inline]
625    pub fn includes_orca_whirlpool(&self) -> bool {
626        self.includes_any(&[
627            EventType::OrcaWhirlpoolSwap,
628            EventType::OrcaWhirlpoolLiquidityIncreased,
629            EventType::OrcaWhirlpoolLiquidityDecreased,
630            EventType::OrcaWhirlpoolPoolInitialized,
631        ])
632    }
633
634    #[inline]
635    pub fn includes_meteora_pools(&self) -> bool {
636        self.includes_any(&[
637            EventType::MeteoraPoolsSwap,
638            EventType::MeteoraPoolsAddLiquidity,
639            EventType::MeteoraPoolsRemoveLiquidity,
640            EventType::MeteoraPoolsBootstrapLiquidity,
641            EventType::MeteoraPoolsPoolCreated,
642            EventType::MeteoraPoolsSetPoolFees,
643        ])
644    }
645
646    #[inline]
647    pub fn includes_meteora_dlmm(&self) -> bool {
648        self.includes_any(&[
649            EventType::MeteoraDlmmSwap,
650            EventType::MeteoraDlmmAddLiquidity,
651            EventType::MeteoraDlmmRemoveLiquidity,
652            EventType::MeteoraDlmmInitializePool,
653            EventType::MeteoraDlmmInitializeBinArray,
654            EventType::MeteoraDlmmCreatePosition,
655            EventType::MeteoraDlmmClosePosition,
656            EventType::MeteoraDlmmClaimFee,
657        ])
658    }
659
660    #[inline]
661    pub fn includes_meteora_dbc(&self) -> bool {
662        self.includes_any(&[
663            EventType::MeteoraDbcSwap,
664            EventType::MeteoraDbcInitializePool,
665            EventType::MeteoraDbcCurveComplete,
666        ])
667    }
668}
669
670#[inline]
671fn pumpfun_trade_filter_is_generic(include_only: &[EventType]) -> bool {
672    include_only.contains(&EventType::PumpFunTrade)
673        && !include_only.iter().any(|t| {
674            matches!(
675                t,
676                EventType::PumpFunBuy | EventType::PumpFunSell | EventType::PumpFunBuyExactSolIn
677            )
678        })
679}
680
681#[inline]
682fn pumpfun_buy_filter_is_generic(include_only: &[EventType]) -> bool {
683    include_only.contains(&EventType::PumpFunBuy)
684        && !include_only.contains(&EventType::PumpFunBuyExactSolIn)
685}
686
687#[inline]
688fn is_pumpfun_create_family(event_type: EventType) -> bool {
689    matches!(event_type, EventType::PumpFunCreate | EventType::PumpFunCreateV2)
690}
691
692#[inline]
693pub fn event_type_from_dex_event(event: &crate::core::events::DexEvent) -> Option<EventType> {
694    use crate::core::events::DexEvent;
695    match event {
696        DexEvent::PumpFunCreate(_) => Some(EventType::PumpFunCreate),
697        DexEvent::PumpFunCreateV2(_) => Some(EventType::PumpFunCreateV2),
698        DexEvent::PumpFunTrade(_) => Some(EventType::PumpFunTrade),
699        DexEvent::PumpFunBuy(_) => Some(EventType::PumpFunBuy),
700        DexEvent::PumpFunSell(_) => Some(EventType::PumpFunSell),
701        DexEvent::PumpFunBuyExactSolIn(_) => Some(EventType::PumpFunBuyExactSolIn),
702        DexEvent::PumpFunMigrate(_) => Some(EventType::PumpFunMigrate),
703        DexEvent::PumpFeesCreateFeeSharingConfig(_) => {
704            Some(EventType::PumpFeesCreateFeeSharingConfig)
705        }
706        DexEvent::PumpFeesInitializeFeeConfig(_) => Some(EventType::PumpFeesInitializeFeeConfig),
707        DexEvent::PumpFeesResetFeeSharingConfig(_) => {
708            Some(EventType::PumpFeesResetFeeSharingConfig)
709        }
710        DexEvent::PumpFeesRevokeFeeSharingAuthority(_) => {
711            Some(EventType::PumpFeesRevokeFeeSharingAuthority)
712        }
713        DexEvent::PumpFeesTransferFeeSharingAuthority(_) => {
714            Some(EventType::PumpFeesTransferFeeSharingAuthority)
715        }
716        DexEvent::PumpFeesUpdateAdmin(_) => Some(EventType::PumpFeesUpdateAdmin),
717        DexEvent::PumpFeesUpdateFeeConfig(_) => Some(EventType::PumpFeesUpdateFeeConfig),
718        DexEvent::PumpFeesUpdateFeeShares(_) => Some(EventType::PumpFeesUpdateFeeShares),
719        DexEvent::PumpFeesUpsertFeeTiers(_) => Some(EventType::PumpFeesUpsertFeeTiers),
720        DexEvent::PumpFunMigrateBondingCurveCreator(_) => {
721            Some(EventType::PumpFunMigrateBondingCurveCreator)
722        }
723        DexEvent::PumpFunGlobalAccount(_) => Some(EventType::AccountPumpFunGlobal),
724        DexEvent::PumpFunBondingCurveAccount(_) => Some(EventType::AccountPumpFunBondingCurve),
725        DexEvent::PumpFunFeeConfigAccount(_) => Some(EventType::AccountPumpFunFeeConfig),
726        DexEvent::PumpFunSharingConfigAccount(_) => Some(EventType::AccountPumpFunSharingConfig),
727        DexEvent::PumpFunGlobalVolumeAccumulatorAccount(_) => {
728            Some(EventType::AccountPumpFunGlobalVolumeAccumulator)
729        }
730        DexEvent::PumpFunUserVolumeAccumulatorAccount(_) => {
731            Some(EventType::AccountPumpFunUserVolumeAccumulator)
732        }
733        DexEvent::PumpSwapTrade(_) => Some(EventType::PumpSwapTrade),
734        DexEvent::PumpSwapBuy(_) => Some(EventType::PumpSwapBuy),
735        DexEvent::PumpSwapSell(_) => Some(EventType::PumpSwapSell),
736        DexEvent::PumpSwapCreatePool(_) => Some(EventType::PumpSwapCreatePool),
737        DexEvent::PumpSwapLiquidityAdded(_) => Some(EventType::PumpSwapLiquidityAdded),
738        DexEvent::PumpSwapLiquidityRemoved(_) => Some(EventType::PumpSwapLiquidityRemoved),
739        DexEvent::MeteoraDammV2Swap(_) => Some(EventType::MeteoraDammV2Swap),
740        DexEvent::MeteoraDammV2CreatePosition(_) => Some(EventType::MeteoraDammV2CreatePosition),
741        DexEvent::MeteoraDammV2ClosePosition(_) => Some(EventType::MeteoraDammV2ClosePosition),
742        DexEvent::MeteoraDammV2AddLiquidity(_) => Some(EventType::MeteoraDammV2AddLiquidity),
743        DexEvent::MeteoraDammV2RemoveLiquidity(_) => Some(EventType::MeteoraDammV2RemoveLiquidity),
744        DexEvent::MeteoraDammV2InitializePool(_) => Some(EventType::MeteoraDammV2InitializePool),
745        DexEvent::MeteoraDammV2UpdateDelegatePermission(_) => {
746            Some(EventType::MeteoraDammV2UpdateDelegatePermission)
747        }
748        DexEvent::MeteoraDammV2WithdrawDeadLiquidityReward(_) => {
749            Some(EventType::MeteoraDammV2WithdrawDeadLiquidityReward)
750        }
751        DexEvent::MeteoraDammV2CreateConfig(_) => Some(EventType::MeteoraDammV2CreateConfig),
752        DexEvent::MeteoraDammV2CreateDynamicConfig(_) => {
753            Some(EventType::MeteoraDammV2CreateDynamicConfig)
754        }
755        DexEvent::MeteoraDbcSwap(_) => Some(EventType::MeteoraDbcSwap),
756        DexEvent::MeteoraDbcInitializePool(_) => Some(EventType::MeteoraDbcInitializePool),
757        DexEvent::MeteoraDbcCurveComplete(_) => Some(EventType::MeteoraDbcCurveComplete),
758        DexEvent::RaydiumLaunchlabTrade(_) => Some(EventType::RaydiumLaunchlabTrade),
759        DexEvent::RaydiumLaunchlabPoolCreate(_) => Some(EventType::RaydiumLaunchlabPoolCreate),
760        DexEvent::RaydiumLaunchlabMigrateAmm(_) => Some(EventType::RaydiumLaunchlabMigrateAmm),
761        DexEvent::RaydiumClmmSwap(_) => Some(EventType::RaydiumClmmSwap),
762        DexEvent::RaydiumClmmCreatePool(_) => Some(EventType::RaydiumClmmCreatePool),
763        DexEvent::RaydiumClmmOpenPosition(_) => Some(EventType::RaydiumClmmOpenPosition),
764        DexEvent::RaydiumClmmOpenPositionWithTokenExtNft(_) => {
765            Some(EventType::RaydiumClmmOpenPositionWithTokenExtNft)
766        }
767        DexEvent::RaydiumClmmClosePosition(_) => Some(EventType::RaydiumClmmClosePosition),
768        DexEvent::RaydiumClmmIncreaseLiquidity(_) => Some(EventType::RaydiumClmmIncreaseLiquidity),
769        DexEvent::RaydiumClmmDecreaseLiquidity(_) => Some(EventType::RaydiumClmmDecreaseLiquidity),
770        DexEvent::RaydiumClmmLiquidityChange(_) => Some(EventType::RaydiumClmmLiquidityChange),
771        DexEvent::RaydiumClmmConfigChange(_) => Some(EventType::RaydiumClmmConfigChange),
772        DexEvent::RaydiumClmmCreatePersonalPosition(_) => {
773            Some(EventType::RaydiumClmmCreatePersonalPosition)
774        }
775        DexEvent::RaydiumClmmLiquidityCalculate(_) => {
776            Some(EventType::RaydiumClmmLiquidityCalculate)
777        }
778        DexEvent::RaydiumClmmOpenLimitOrder(_) => Some(EventType::RaydiumClmmOpenLimitOrder),
779        DexEvent::RaydiumClmmIncreaseLimitOrder(_) => {
780            Some(EventType::RaydiumClmmIncreaseLimitOrder)
781        }
782        DexEvent::RaydiumClmmDecreaseLimitOrder(_) => {
783            Some(EventType::RaydiumClmmDecreaseLimitOrder)
784        }
785        DexEvent::RaydiumClmmSettleLimitOrder(_) => Some(EventType::RaydiumClmmSettleLimitOrder),
786        DexEvent::RaydiumClmmUpdateRewardInfos(_) => Some(EventType::RaydiumClmmUpdateRewardInfos),
787        DexEvent::RaydiumClmmCollectFee(_) => Some(EventType::RaydiumClmmCollectFee),
788        DexEvent::RaydiumClmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumClmmAmmConfig),
789        DexEvent::RaydiumClmmPoolStateAccount(_) => Some(EventType::AccountRaydiumClmmPoolState),
790        DexEvent::RaydiumClmmTickArrayStateAccount(_) => {
791            Some(EventType::AccountRaydiumClmmTickArrayState)
792        }
793        DexEvent::RaydiumCpmmSwap(_) => Some(EventType::RaydiumCpmmSwap),
794        DexEvent::RaydiumCpmmDeposit(_) => Some(EventType::RaydiumCpmmDeposit),
795        DexEvent::RaydiumCpmmWithdraw(_) => Some(EventType::RaydiumCpmmWithdraw),
796        DexEvent::RaydiumCpmmInitialize(_) => Some(EventType::RaydiumCpmmInitialize),
797        DexEvent::RaydiumCpmmAmmConfigAccount(_) => Some(EventType::AccountRaydiumCpmmAmmConfig),
798        DexEvent::RaydiumCpmmPoolStateAccount(_) => Some(EventType::AccountRaydiumCpmmPoolState),
799        DexEvent::RaydiumAmmV4Swap(_) => Some(EventType::RaydiumAmmV4Swap),
800        DexEvent::RaydiumAmmV4Deposit(_) => Some(EventType::RaydiumAmmV4Deposit),
801        DexEvent::RaydiumAmmV4Initialize2(_) => Some(EventType::RaydiumAmmV4Initialize2),
802        DexEvent::RaydiumAmmV4Withdraw(_) => Some(EventType::RaydiumAmmV4Withdraw),
803        DexEvent::RaydiumAmmV4WithdrawPnl(_) => Some(EventType::RaydiumAmmV4WithdrawPnl),
804        DexEvent::OrcaWhirlpoolSwap(_) => Some(EventType::OrcaWhirlpoolSwap),
805        DexEvent::OrcaWhirlpoolLiquidityIncreased(_) => {
806            Some(EventType::OrcaWhirlpoolLiquidityIncreased)
807        }
808        DexEvent::OrcaWhirlpoolLiquidityDecreased(_) => {
809            Some(EventType::OrcaWhirlpoolLiquidityDecreased)
810        }
811        DexEvent::OrcaWhirlpoolPoolInitialized(_) => Some(EventType::OrcaWhirlpoolPoolInitialized),
812        DexEvent::OrcaWhirlpoolAccount(_) => Some(EventType::AccountOrcaWhirlpool),
813        DexEvent::OrcaPositionAccount(_) => Some(EventType::AccountOrcaPosition),
814        DexEvent::OrcaTickArrayAccount(_) => Some(EventType::AccountOrcaTickArray),
815        DexEvent::LiquidityAccountSnapshot(_) => Some(EventType::AccountLiquiditySnapshot),
816        DexEvent::RawAccountSnapshot(_) => Some(EventType::AccountRawSnapshot),
817        DexEvent::OrcaFeeTierAccount(_) => Some(EventType::AccountOrcaFeeTier),
818        DexEvent::OrcaWhirlpoolsConfigAccount(_) => Some(EventType::AccountOrcaWhirlpoolsConfig),
819        DexEvent::MeteoraPoolsSwap(_) => Some(EventType::MeteoraPoolsSwap),
820        DexEvent::MeteoraPoolsAddLiquidity(_) => Some(EventType::MeteoraPoolsAddLiquidity),
821        DexEvent::MeteoraPoolsRemoveLiquidity(_) => Some(EventType::MeteoraPoolsRemoveLiquidity),
822        DexEvent::MeteoraPoolsBootstrapLiquidity(_) => {
823            Some(EventType::MeteoraPoolsBootstrapLiquidity)
824        }
825        DexEvent::MeteoraPoolsPoolCreated(_) => Some(EventType::MeteoraPoolsPoolCreated),
826        DexEvent::MeteoraPoolsSetPoolFees(_) => Some(EventType::MeteoraPoolsSetPoolFees),
827        DexEvent::MeteoraDlmmSwap(_) => Some(EventType::MeteoraDlmmSwap),
828        DexEvent::MeteoraDlmmAddLiquidity(_) => Some(EventType::MeteoraDlmmAddLiquidity),
829        DexEvent::MeteoraDlmmRemoveLiquidity(_) => Some(EventType::MeteoraDlmmRemoveLiquidity),
830        DexEvent::MeteoraDlmmInitializePool(_) => Some(EventType::MeteoraDlmmInitializePool),
831        DexEvent::MeteoraDlmmInitializeBinArray(_) => {
832            Some(EventType::MeteoraDlmmInitializeBinArray)
833        }
834        DexEvent::MeteoraDlmmCreatePosition(_) => Some(EventType::MeteoraDlmmCreatePosition),
835        DexEvent::MeteoraDlmmClosePosition(_) => Some(EventType::MeteoraDlmmClosePosition),
836        DexEvent::MeteoraDlmmClaimFee(_) => Some(EventType::MeteoraDlmmClaimFee),
837        DexEvent::TokenAccount(_) => Some(EventType::TokenAccount),
838        DexEvent::TokenInfo(_) => Some(EventType::TokenInfo),
839        DexEvent::NonceAccount(_) => Some(EventType::NonceAccount),
840        DexEvent::PumpSwapGlobalConfigAccount(_) => Some(EventType::AccountPumpSwapGlobalConfig),
841        DexEvent::PumpSwapPoolAccount(_) => Some(EventType::AccountPumpSwapPool),
842        DexEvent::BlockMeta(_) => Some(EventType::BlockMeta),
843        DexEvent::Error(_) => None,
844    }
845}
846
847#[cfg(test)]
848mod event_type_filter_tests {
849    use super::*;
850
851    #[test]
852    fn generic_trade_filters_cover_specific_trade_variants() {
853        let pump = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
854        assert!(pump.should_include(EventType::PumpFunTrade));
855        assert!(pump.should_include(EventType::PumpFunBuy));
856        assert!(pump.should_include(EventType::PumpFunSell));
857        assert!(pump.should_include(EventType::PumpFunBuyExactSolIn));
858
859        let pump_specific = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
860        assert!(!pump_specific.should_include(EventType::PumpFunTrade));
861        assert!(pump_specific.should_include(EventType::PumpFunBuy));
862        assert!(!pump_specific.should_include(EventType::PumpFunSell));
863        assert!(pump_specific.should_include(EventType::PumpFunBuyExactSolIn));
864
865        let pump_exact_buy = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
866        assert!(!pump_exact_buy.should_include(EventType::PumpFunTrade));
867        assert!(!pump_exact_buy.should_include(EventType::PumpFunBuy));
868        assert!(!pump_exact_buy.should_include(EventType::PumpFunSell));
869        assert!(pump_exact_buy.should_include(EventType::PumpFunBuyExactSolIn));
870
871        let pumpswap = EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]);
872        assert!(pumpswap.should_include(EventType::PumpSwapBuy));
873        assert!(pumpswap.should_include(EventType::PumpSwapSell));
874
875        let exclude_pumpswap = EventTypeFilter::exclude_types(vec![EventType::PumpSwapTrade]);
876        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapBuy));
877        assert!(!exclude_pumpswap.should_include(EventType::PumpSwapSell));
878    }
879
880    #[test]
881    fn generic_pumpfun_trade_filter_normalizes_specific_variants() {
882        use crate::core::events::{DexEvent, PumpFunTradeEvent};
883
884        let filter = EventTypeFilter::include_only(vec![EventType::PumpFunTrade]);
885        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
886        assert!(matches!(filter.normalize_dex_event(event), DexEvent::PumpFunTrade(_)));
887
888        let specific_filter =
889            EventTypeFilter::include_only(vec![EventType::PumpFunTrade, EventType::PumpFunBuy]);
890        let event = DexEvent::PumpFunBuy(PumpFunTradeEvent { is_buy: true, ..Default::default() });
891        assert!(matches!(specific_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
892
893        let buy_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuy]);
894        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
895            is_buy: true,
896            ..Default::default()
897        });
898        assert!(matches!(buy_filter.normalize_dex_event(event), DexEvent::PumpFunBuy(_)));
899
900        let exact_filter = EventTypeFilter::include_only(vec![EventType::PumpFunBuyExactSolIn]);
901        let event = DexEvent::PumpFunBuyExactSolIn(PumpFunTradeEvent {
902            is_buy: true,
903            ..Default::default()
904        });
905        assert!(matches!(
906            exact_filter.normalize_dex_event(event),
907            DexEvent::PumpFunBuyExactSolIn(_)
908        ));
909
910        let create_and_trade_filter =
911            EventTypeFilter::include_only(vec![EventType::PumpFunCreate, EventType::PumpFunTrade]);
912        let event =
913            DexEvent::PumpFunSell(PumpFunTradeEvent { is_buy: false, ..Default::default() });
914        assert!(matches!(
915            create_and_trade_filter.normalize_dex_event(event),
916            DexEvent::PumpFunTrade(_)
917        ));
918    }
919
920    #[test]
921    fn all_protocol_groups_are_filterable() {
922        assert!(EventTypeFilter::include_only(vec![EventType::PumpFunTrade]).includes_pumpfun());
923        assert!(!EventTypeFilter::include_only(vec![EventType::AccountPumpFunGlobal])
924            .includes_pumpfun());
925        assert!(
926            !EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateAdmin]).includes_pumpfun()
927        );
928        assert!(EventTypeFilter::include_only(vec![EventType::PumpSwapTrade]).includes_pumpswap());
929        assert!(EventTypeFilter::include_only(vec![EventType::PumpFeesUpdateFeeShares])
930            .includes_pump_fees());
931        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumLaunchlabTrade])
932            .includes_raydium_launchlab());
933        assert!(
934            EventTypeFilter::include_only(vec![EventType::RaydiumCpmmSwap]).includes_raydium_cpmm()
935        );
936        assert!(
937            EventTypeFilter::include_only(vec![EventType::RaydiumClmmSwap]).includes_raydium_clmm()
938        );
939        assert!(!EventTypeFilter::include_only(vec![EventType::AccountRaydiumClmmPoolState])
940            .includes_raydium_clmm());
941        assert!(EventTypeFilter::include_only(vec![EventType::RaydiumAmmV4Swap])
942            .includes_raydium_amm_v4());
943        assert!(EventTypeFilter::include_only(vec![EventType::OrcaWhirlpoolSwap])
944            .includes_orca_whirlpool());
945        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraPoolsSwap])
946            .includes_meteora_pools());
947        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2Swap])
948            .includes_meteora_damm_v2());
949        assert!(EventTypeFilter::include_only(vec![EventType::MeteoraDammV2InitializePool])
950            .includes_meteora_damm_v2());
951        assert!(
952            EventTypeFilter::include_only(vec![EventType::MeteoraDlmmSwap]).includes_meteora_dlmm()
953        );
954        assert!(
955            EventTypeFilter::include_only(vec![EventType::MeteoraDbcSwap]).includes_meteora_dbc()
956        );
957    }
958
959    #[test]
960    fn exclude_filters_do_not_skip_whole_protocol_groups() {
961        let raydium = EventTypeFilter::exclude_types(vec![EventType::RaydiumCpmmSwap]);
962        assert!(raydium.includes_raydium_cpmm());
963        assert!(!raydium.should_include(EventType::RaydiumCpmmSwap));
964        assert!(raydium.should_include(EventType::RaydiumCpmmDeposit));
965
966        let all_cpmm = EventTypeFilter::exclude_types(vec![
967            EventType::RaydiumCpmmSwap,
968            EventType::RaydiumCpmmDeposit,
969            EventType::RaydiumCpmmWithdraw,
970            EventType::RaydiumCpmmInitialize,
971        ]);
972        assert!(!all_cpmm.includes_raydium_cpmm());
973
974        let all_launchlab = EventTypeFilter::exclude_types(vec![
975            EventType::RaydiumLaunchlabTrade,
976            EventType::RaydiumLaunchlabPoolCreate,
977            EventType::RaydiumLaunchlabMigrateAmm,
978        ]);
979        assert!(!all_launchlab.includes_raydium_launchlab());
980
981        let pump = EventTypeFilter::exclude_types(vec![EventType::PumpFunBuy]);
982        assert!(pump.includes_pumpfun());
983        assert!(!pump.should_include(EventType::PumpFunBuy));
984        assert!(!pump.should_include(EventType::PumpFunBuyExactSolIn));
985        assert!(pump.should_include(EventType::PumpFunSell));
986    }
987}
988
989#[derive(Debug, Clone)]
990pub struct SlotFilter {
991    pub min_slot: Option<u64>,
992    pub max_slot: Option<u64>,
993}
994
995impl SlotFilter {
996    pub fn new() -> Self {
997        Self { min_slot: None, max_slot: None }
998    }
999
1000    pub fn min_slot(mut self, slot: u64) -> Self {
1001        self.min_slot = Some(slot);
1002        self
1003    }
1004
1005    pub fn max_slot(mut self, slot: u64) -> Self {
1006        self.max_slot = Some(slot);
1007        self
1008    }
1009}
1010
1011impl Default for SlotFilter {
1012    fn default() -> Self {
1013        Self::new()
1014    }
1015}