diff --git a/src/config_loader.py b/src/config_loader.py index a9191e3..7feb1cb 100644 --- a/src/config_loader.py +++ b/src/config_loader.py @@ -45,7 +45,7 @@ VALID_VALUES = { # Platform-specific listener compatibility PLATFORM_LISTENER_COMPATIBILITY = { Platform.PUMP_FUN: ["logs", "blocks", "geyser", "pumpportal"], - Platform.LETS_BONK: ["blocks", "geyser"], # LetsBonk only supports geyser and blocks + Platform.LETS_BONK: ["blocks", "geyser", "pumpportal"], } diff --git a/src/monitoring/listener_factory.py b/src/monitoring/listener_factory.py index e10b5d8..793bbc1 100644 --- a/src/monitoring/listener_factory.py +++ b/src/monitoring/listener_factory.py @@ -87,21 +87,30 @@ class ListenerFactory: return listener elif listener_type == "pumpportal": - # PumpPortal is pump.fun specific, so filter platforms - pumpfun_platforms = [Platform.PUMP_FUN] - if platforms: - pumpfun_platforms = [p for p in platforms if p == Platform.PUMP_FUN] - - if not pumpfun_platforms: - raise ValueError("PumpPortal listener only supports pump.fun platform") - - from monitoring.pumpportal_listener import PumpPortalListener - - listener = PumpPortalListener( - pump_program=None, # Will be determined from platform - pumpportal_url=pumpportal_url, + # Import the new universal PumpPortal listener + from monitoring.universal_pumpportal_listener import ( + UniversalPumpPortalListener, ) - logger.info("Created PumpPortal listener for token monitoring") + + # Validate that requested platforms support PumpPortal + supported_pumpportal_platforms = [Platform.PUMP_FUN, Platform.LETS_BONK] + + if platforms: + unsupported = [p for p in platforms if p not in supported_pumpportal_platforms] + if unsupported: + logger.warning(f"Platforms {[p.value for p in unsupported]} do not support PumpPortal") + + # Filter to only supported platforms + filtered_platforms = [p for p in platforms if p in supported_pumpportal_platforms] + if not filtered_platforms: + raise ValueError("No supported platforms specified for PumpPortal listener") + platforms = filtered_platforms + + listener = UniversalPumpPortalListener( + pumpportal_url=pumpportal_url, + platforms=platforms, + ) + logger.info(f"Created Universal PumpPortal listener for platforms: {[p.value for p in (platforms or supported_pumpportal_platforms)]}") return listener else: @@ -132,6 +141,15 @@ class ListenerFactory: if platform == Platform.PUMP_FUN: return ["logs", "blocks", "geyser", "pumpportal"] elif platform == Platform.LETS_BONK: - return ["blocks", "geyser"] # LetsBonk only supports geyser and blocks + return ["blocks", "geyser", "pumpportal"] # Added pumpportal support else: - return ["blocks", "geyser"] # Default universal listeners \ No newline at end of file + return ["blocks", "geyser"] # Default universal listeners + + @staticmethod + def get_pumpportal_supported_platforms() -> list[Platform]: + """Get list of platforms that support PumpPortal listener. + + Returns: + List of platforms with PumpPortal support + """ + return [Platform.PUMP_FUN, Platform.LETS_BONK] \ No newline at end of file diff --git a/src/monitoring/pumpportal_listener.py b/src/monitoring/universal_pumpportal_listener.py similarity index 60% rename from src/monitoring/pumpportal_listener.py rename to src/monitoring/universal_pumpportal_listener.py index 6f8ab63..a3778b5 100644 --- a/src/monitoring/pumpportal_listener.py +++ b/src/monitoring/universal_pumpportal_listener.py @@ -1,5 +1,5 @@ """ -PumpPortal monitoring for pump.fun tokens. +Universal PumpPortal listener that works with multiple platforms. """ import asyncio @@ -7,36 +7,58 @@ import json from collections.abc import Awaitable, Callable import websockets -from solders.pubkey import Pubkey -from interfaces.core import TokenInfo +from interfaces.core import Platform, TokenInfo from monitoring.base_listener import BaseTokenListener -from monitoring.pumpportal_event_processor import PumpPortalEventProcessor -from platforms.pumpfun.address_provider import PumpFunAddresses from utils.logger import get_logger logger = get_logger(__name__) -class PumpPortalListener(BaseTokenListener): - """PumpPortal listener for pump.fun token creation events.""" +class UniversalPumpPortalListener(BaseTokenListener): + """Universal PumpPortal listener that works with multiple platforms.""" def __init__( self, - pump_program: Pubkey | None = None, pumpportal_url: str = "wss://pumpportal.fun/api/data", + platforms: list[Platform] | None = None, ): - """Initialize token listener. + """Initialize universal PumpPortal listener. Args: - pump_program: Pump.fun program address (optional, will use default if None) pumpportal_url: PumpPortal WebSocket URL + platforms: List of platforms to monitor (if None, monitor all supported platforms) """ super().__init__() - self.pump_program = pump_program or PumpFunAddresses.PROGRAM self.pumpportal_url = pumpportal_url - self.event_processor = PumpPortalEventProcessor(self.pump_program) self.ping_interval = 20 # seconds + + # Get platform-specific processors + from platforms.letsbonk.pumpportal_processor import LetsBonkPumpPortalProcessor + from platforms.pumpfun.pumpportal_processor import PumpFunPumpPortalProcessor + + # Create processor instances + all_processors = [ + PumpFunPumpPortalProcessor(), + LetsBonkPumpPortalProcessor(), + ] + + # Filter processors based on requested platforms + if platforms is None: + self.processors = all_processors + else: + self.processors = [p for p in all_processors if p.platform in platforms] + + # Build mapping of pool names to processors for quick lookup + self.pool_to_processors: dict[str, list] = {} + for processor in self.processors: + for pool_name in processor.supported_pool_names: + if pool_name not in self.pool_to_processors: + self.pool_to_processors[pool_name] = [] + self.pool_to_processors[pool_name].append(processor) + + logger.info(f"Initialized Universal PumpPortal listener for platforms: {[p.platform.value for p in self.processors]}") + logger.info(f"Monitoring pools: {list(self.pool_to_processors.keys())}") async def listen_for_tokens( self, @@ -64,9 +86,10 @@ class PumpPortalListener(BaseTokenListener): continue logger.info( - f"New token detected: {token_info.name} ({token_info.symbol})" + f"New token detected: {token_info.name} ({token_info.symbol}) on {token_info.platform.value}" ) + # Apply filters if match_string and not ( match_string.lower() in token_info.name.lower() or match_string.lower() in token_info.symbol.lower() @@ -76,14 +99,14 @@ class PumpPortalListener(BaseTokenListener): ) continue - if ( - creator_address - and str(token_info.user) != creator_address - ): - logger.info( - f"Token not created by {creator_address}. Skipping..." - ) - continue + if creator_address: + creator_str = str(token_info.creator) if token_info.creator else "" + user_str = str(token_info.user) if token_info.user else "" + if creator_address not in [creator_str, user_str]: + logger.info( + f"Token not created by {creator_address}. Skipping..." + ) + continue await token_callback(token_info) @@ -150,18 +173,35 @@ class PumpPortalListener(BaseTokenListener): data = json.loads(response) # Handle different message formats from PumpPortal - token_info = None + token_data = None if "method" in data and data["method"] == "newToken": # Standard newToken method format params = data.get("params", []) if params and len(params) > 0: token_data = params[0] - token_info = self.event_processor.process_token_data(token_data) - elif "signature" in data and "mint" in data: + elif "signature" in data and "mint" in data and "pool" in data: # Direct token data format - token_info = self.event_processor.process_token_data(data) + token_data = data - return token_info + if not token_data: + return None + + # Get pool name to determine which processor to use + pool_name = token_data.get("pool", "").lower() + if pool_name not in self.pool_to_processors: + logger.debug(f"Ignoring token from unsupported pool: {pool_name}") + return None + + # Try each processor that supports this pool + for processor in self.pool_to_processors[pool_name]: + if processor.can_process(token_data): + token_info = processor.process_token_data(token_data) + if token_info: + logger.debug(f"Successfully processed token using {processor.platform.value} processor") + return token_info + + logger.debug(f"No processor could handle token data from pool {pool_name}") + return None except TimeoutError: logger.debug("No data received from PumpPortal for 30 seconds") diff --git a/src/platforms/letsbonk/__init__.py b/src/platforms/letsbonk/__init__.py index 60c1cfc..110962b 100644 --- a/src/platforms/letsbonk/__init__.py +++ b/src/platforms/letsbonk/__init__.py @@ -9,11 +9,13 @@ from .address_provider import LetsBonkAddressProvider from .curve_manager import LetsBonkCurveManager from .event_parser import LetsBonkEventParser from .instruction_builder import LetsBonkInstructionBuilder +from .pumpportal_processor import LetsBonkPumpPortalProcessor # Export implementations for direct use if needed __all__ = [ 'LetsBonkAddressProvider', 'LetsBonkCurveManager', 'LetsBonkEventParser', - 'LetsBonkInstructionBuilder' + 'LetsBonkInstructionBuilder', + 'LetsBonkPumpPortalProcessor' ] \ No newline at end of file diff --git a/src/platforms/letsbonk/address_provider.py b/src/platforms/letsbonk/address_provider.py index 647f093..37d392d 100644 --- a/src/platforms/letsbonk/address_provider.py +++ b/src/platforms/letsbonk/address_provider.py @@ -79,6 +79,52 @@ class LetsBonkAddressProvider(AddressProvider): ) return pool_state + def derive_base_vault(self, base_mint: Pubkey, quote_mint: Pubkey | None = None) -> Pubkey: + """Derive the base vault address for a token pair. + + Args: + base_mint: Base token mint address + quote_mint: Quote token mint (defaults to WSOL) + + Returns: + Base vault address + """ + if quote_mint is None: + quote_mint = SystemAddresses.SOL_MINT + + # First derive the pool state address + pool_state = self.derive_pool_address(base_mint, quote_mint) + + # Then derive the base vault using pool_vault seed + base_vault, _ = Pubkey.find_program_address( + [b"pool_vault", bytes(pool_state), bytes(base_mint)], + LetsBonkAddresses.PROGRAM + ) + return base_vault + + def derive_quote_vault(self, base_mint: Pubkey, quote_mint: Pubkey | None = None) -> Pubkey: + """Derive the quote vault address for a token pair. + + Args: + base_mint: Base token mint address + quote_mint: Quote token mint (defaults to WSOL) + + Returns: + Quote vault address + """ + if quote_mint is None: + quote_mint = SystemAddresses.SOL_MINT + + # First derive the pool state address + pool_state = self.derive_pool_address(base_mint, quote_mint) + + # Then derive the quote vault using pool_vault seed + quote_vault, _ = Pubkey.find_program_address( + [b"pool_vault", bytes(pool_state), bytes(quote_mint)], + LetsBonkAddresses.PROGRAM + ) + return quote_vault + def derive_user_token_account(self, user: Pubkey, mint: Pubkey) -> Pubkey: """Derive user's associated token account address. @@ -102,15 +148,22 @@ class LetsBonkAddressProvider(AddressProvider): """ accounts = {} - # Add pool state - must be present in token_info + # Add pool state - derive if not present or use existing if token_info.pool_state: accounts["pool_state"] = token_info.pool_state + else: + accounts["pool_state"] = self.derive_pool_address(token_info.mint) - # Add vault addresses - must be present in token_info + # Add vault addresses - derive if not present or use existing if token_info.base_vault: accounts["base_vault"] = token_info.base_vault + else: + accounts["base_vault"] = self.derive_base_vault(token_info.mint) + if token_info.quote_vault: accounts["quote_vault"] = token_info.quote_vault + else: + accounts["quote_vault"] = self.derive_quote_vault(token_info.mint) # Derive authority PDA accounts["authority"] = self.derive_authority_pda() @@ -174,20 +227,15 @@ class LetsBonkAddressProvider(AddressProvider): """ additional_accounts = self.get_additional_accounts(token_info) - # Vault addresses must be present in token_info - if not token_info.base_vault or not token_info.quote_vault: - raise ValueError(f"Missing required vault addresses for token {token_info.mint}. " - f"base_vault: {token_info.base_vault}, quote_vault: {token_info.quote_vault}") - return { "payer": user, "authority": additional_accounts["authority"], "global_config": LetsBonkAddresses.GLOBAL_CONFIG, "platform_config": LetsBonkAddresses.PLATFORM_CONFIG, - "pool_state": token_info.pool_state, + "pool_state": additional_accounts["pool_state"], "user_base_token": self.derive_user_token_account(user, token_info.mint), - "base_vault": token_info.base_vault, - "quote_vault": token_info.quote_vault, + "base_vault": additional_accounts["base_vault"], + "quote_vault": additional_accounts["quote_vault"], "base_token_mint": token_info.mint, "quote_token_mint": SystemAddresses.SOL_MINT, "base_token_program": SystemAddresses.TOKEN_PROGRAM, @@ -208,20 +256,15 @@ class LetsBonkAddressProvider(AddressProvider): """ additional_accounts = self.get_additional_accounts(token_info) - # Vault addresses must be present in token_info - if not token_info.base_vault or not token_info.quote_vault: - raise ValueError(f"Missing required vault addresses for token {token_info.mint}. " - f"base_vault: {token_info.base_vault}, quote_vault: {token_info.quote_vault}") - return { "payer": user, "authority": additional_accounts["authority"], "global_config": LetsBonkAddresses.GLOBAL_CONFIG, "platform_config": LetsBonkAddresses.PLATFORM_CONFIG, - "pool_state": token_info.pool_state, + "pool_state": additional_accounts["pool_state"], "user_base_token": self.derive_user_token_account(user, token_info.mint), - "base_vault": token_info.base_vault, - "quote_vault": token_info.quote_vault, + "base_vault": additional_accounts["base_vault"], + "quote_vault": additional_accounts["quote_vault"], "base_token_mint": token_info.mint, "quote_token_mint": SystemAddresses.SOL_MINT, "base_token_program": SystemAddresses.TOKEN_PROGRAM, diff --git a/src/platforms/letsbonk/pumpportal_processor.py b/src/platforms/letsbonk/pumpportal_processor.py new file mode 100644 index 0000000..4af1431 --- /dev/null +++ b/src/platforms/letsbonk/pumpportal_processor.py @@ -0,0 +1,118 @@ +""" +LetsBonk-specific PumpPortal event processor. +File: src/platforms/letsbonk/pumpportal_processor.py +""" + +from solders.pubkey import Pubkey + +from interfaces.core import Platform, TokenInfo +from platforms.letsbonk.address_provider import LetsBonkAddressProvider +from utils.logger import get_logger + +logger = get_logger(__name__) + + +class LetsBonkPumpPortalProcessor: + """PumpPortal processor for LetsBonk tokens.""" + + def __init__(self): + """Initialize the processor with address provider.""" + self.address_provider = LetsBonkAddressProvider() + + @property + def platform(self) -> Platform: + """Get the platform this processor handles.""" + return Platform.LETS_BONK + + @property + def supported_pool_names(self) -> list[str]: + """Get the pool names this processor supports from PumpPortal.""" + return ["bonk"] # PumpPortal pool name for LetsBonk/bonk pools + + def can_process(self, token_data: dict) -> bool: + """Check if this processor can handle the given token data. + + Args: + token_data: Token data from PumpPortal + + Returns: + True if this processor can handle the token data + """ + pool = token_data.get("pool", "").lower() + return pool in self.supported_pool_names + + def process_token_data(self, token_data: dict) -> TokenInfo | None: + """Process LetsBonk token data from PumpPortal. + + Args: + token_data: Token data from PumpPortal WebSocket + + Returns: + TokenInfo if token creation found, None otherwise + """ + try: + # Extract required fields for LetsBonk + name = token_data.get("name", "") + symbol = token_data.get("symbol", "") + mint_str = token_data.get("mint") + creator_str = token_data.get("traderPublicKey") + uri = token_data.get("uri", "") + + # Note: LetsBonk tokens from PumpPortal might have different field mappings + # This would need to be adjusted based on actual PumpPortal data for LetsBonk tokens + + if not all([name, symbol, mint_str, creator_str]): + logger.warning("Missing required fields in PumpPortal LetsBonk token data") + return None + + # Convert string addresses to Pubkey objects + mint = Pubkey.from_string(mint_str) + user = Pubkey.from_string(creator_str) + creator = user + + # Derive LetsBonk-specific addresses + pool_state = self.address_provider.derive_pool_address(mint) + + # For LetsBonk, vault addresses might need to be derived differently + # or provided in the PumpPortal data. For now, we'll derive them + # using the standard pattern, but this might need adjustment + additional_accounts = self.address_provider.get_additional_accounts( + # Create a minimal TokenInfo to get additional accounts + TokenInfo( + name=name, + symbol=symbol, + uri=uri, + mint=mint, + platform=Platform.LETS_BONK, + pool_state=pool_state, + user=user, + creator=creator, + base_vault=None, # Will be filled from additional_accounts + quote_vault=None, # Will be filled from additional_accounts + ) + ) + + # Extract vault addresses if available + base_vault = additional_accounts.get("base_vault") + quote_vault = additional_accounts.get("quote_vault") + + # If vaults aren't available from additional_accounts, + # we might need to derive them or leave them None + # and let the trading logic handle the derivation + + return TokenInfo( + name=name, + symbol=symbol, + uri=uri, + mint=mint, + platform=Platform.LETS_BONK, + pool_state=pool_state, + base_vault=base_vault, + quote_vault=quote_vault, + user=user, + creator=creator, + ) + + except Exception as e: + logger.error(f"Failed to process PumpPortal LetsBonk token data: {e}") + return None \ No newline at end of file diff --git a/src/platforms/pumpfun/__init__.py b/src/platforms/pumpfun/__init__.py index bc31eb3..fffc606 100644 --- a/src/platforms/pumpfun/__init__.py +++ b/src/platforms/pumpfun/__init__.py @@ -9,11 +9,13 @@ from .address_provider import PumpFunAddressProvider from .curve_manager import PumpFunCurveManager from .event_parser import PumpFunEventParser from .instruction_builder import PumpFunInstructionBuilder +from .pumpportal_processor import PumpFunPumpPortalProcessor # Export implementations for direct use if needed __all__ = [ 'PumpFunAddressProvider', - 'PumpFunCurveManager', + 'PumpFunCurveManager', 'PumpFunEventParser', - 'PumpFunInstructionBuilder' + 'PumpFunInstructionBuilder', + 'PumpFunPumpPortalProcessor' ] \ No newline at end of file diff --git a/src/monitoring/pumpportal_event_processor.py b/src/platforms/pumpfun/pumpportal_processor.py similarity index 53% rename from src/monitoring/pumpportal_event_processor.py rename to src/platforms/pumpfun/pumpportal_processor.py index 5033e00..dc98d6b 100644 --- a/src/monitoring/pumpportal_event_processor.py +++ b/src/platforms/pumpfun/pumpportal_processor.py @@ -1,34 +1,52 @@ """ -Event processing for pump.fun tokens using PumpPortal data. +PumpFun-specific PumpPortal event processor. +File: src/platforms/pumpfun/pumpportal_processor.py """ from solders.pubkey import Pubkey -from core.pubkeys import SystemAddresses from interfaces.core import Platform, TokenInfo -from platforms.pumpfun.address_provider import PumpFunAddresses +from platforms.pumpfun.address_provider import PumpFunAddressProvider from utils.logger import get_logger logger = get_logger(__name__) -class PumpPortalEventProcessor: - """Processes token creation events from PumpPortal WebSocket.""" - - def __init__(self, pump_program: Pubkey | None = None): - """Initialize event processor. - +class PumpFunPumpPortalProcessor: + """PumpPortal processor for pump.fun tokens.""" + + def __init__(self): + """Initialize the processor with address provider.""" + self.address_provider = PumpFunAddressProvider() + + @property + def platform(self) -> Platform: + """Get the platform this processor handles.""" + return Platform.PUMP_FUN + + @property + def supported_pool_names(self) -> list[str]: + """Get the pool names this processor supports from PumpPortal.""" + return ["pump"] # PumpPortal pool name for pump.fun + + def can_process(self, token_data: dict) -> bool: + """Check if this processor can handle the given token data. + Args: - pump_program: Pump.fun program address (optional, will use default if None) + token_data: Token data from PumpPortal + + Returns: + True if this processor can handle the token data """ - self.pump_program = pump_program or PumpFunAddresses.PROGRAM - + pool = token_data.get("pool", "").lower() + return pool in self.supported_pool_names + def process_token_data(self, token_data: dict) -> TokenInfo | None: - """Process token data from PumpPortal and extract token creation info. - + """Process pump.fun token data from PumpPortal. + Args: token_data: Token data from PumpPortal WebSocket - + Returns: TokenInfo if token creation found, None otherwise """ @@ -41,11 +59,12 @@ class PumpPortalEventProcessor: creator_str = token_data.get("traderPublicKey") # Maps to user field uri = token_data.get("uri", "") - # Additional fields available from PumpPortal but not used: - # - initialBuy: Initial buy amount in SOL - # - marketCapSol: Market cap in SOL + # Additional fields available from PumpPortal but not currently used: + # - initialBuy: Initial buy amount in tokens + # - solAmount: SOL amount spent on initial buy # - vSolInBondingCurve: Virtual SOL in bonding curve # - vTokensInBondingCurve: Virtual tokens in bonding curve + # - marketCapSol: Market cap in SOL # - signature: Transaction signature if not all([name, symbol, mint_str, bonding_curve_str, creator_str]): @@ -61,11 +80,11 @@ class PumpPortalEventProcessor: # since PumpPortal doesn't distinguish between them creator = user - # Calculate derived addresses - associated_bonding_curve = self._find_associated_bonding_curve( + # Derive additional addresses using platform provider + associated_bonding_curve = self.address_provider.derive_associated_bonding_curve( mint, bonding_curve ) - creator_vault = self._find_creator_vault(creator) + creator_vault = self.address_provider.derive_creator_vault(creator) return TokenInfo( name=name, @@ -82,44 +101,4 @@ class PumpPortalEventProcessor: except Exception as e: logger.error(f"Failed to process PumpPortal token data: {e}") - return None - - def _find_associated_bonding_curve( - self, mint: Pubkey, bonding_curve: Pubkey - ) -> Pubkey: - """ - Find the associated bonding curve for a given mint and bonding curve. - This uses the standard ATA derivation. - - Args: - mint: Token mint address - bonding_curve: Bonding curve address - - Returns: - Associated bonding curve address - """ - derived_address, _ = Pubkey.find_program_address( - [ - bytes(bonding_curve), - bytes(SystemAddresses.TOKEN_PROGRAM), - bytes(mint), - ], - SystemAddresses.ASSOCIATED_TOKEN_PROGRAM, - ) - return derived_address - - def _find_creator_vault(self, creator: Pubkey) -> Pubkey: - """ - Find the creator vault for a creator. - - Args: - creator: Creator address - - Returns: - Creator vault address - """ - derived_address, _ = Pubkey.find_program_address( - [b"creator-vault", bytes(creator)], - PumpFunAddresses.PROGRAM, - ) - return derived_address \ No newline at end of file + return None \ No newline at end of file