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