feat(core): pump portal fully integrated

This commit is contained in:
smypmsa
2025-08-04 05:40:02 +00:00
parent bf7fbaa8ab
commit ea1d6b469a
8 changed files with 328 additions and 126 deletions
+34 -16
View File
@@ -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
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]
@@ -1,125 +0,0 @@
"""
Event processing for pump.fun tokens using PumpPortal data.
"""
from solders.pubkey import Pubkey
from core.pubkeys import SystemAddresses
from interfaces.core import Platform, TokenInfo
from platforms.pumpfun.address_provider import PumpFunAddresses
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.
Args:
pump_program: Pump.fun program address (optional, will use default if None)
"""
self.pump_program = pump_program or PumpFunAddresses.PROGRAM
def process_token_data(self, token_data: dict) -> TokenInfo | None:
"""Process token data from PumpPortal and extract token creation info.
Args:
token_data: Token data from PumpPortal WebSocket
Returns:
TokenInfo if token creation found, None otherwise
"""
try:
# Extract required fields
name = token_data.get("name", "")
symbol = token_data.get("symbol", "")
mint_str = token_data.get("mint")
bonding_curve_str = token_data.get("bondingCurveKey")
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
# - vSolInBondingCurve: Virtual SOL in bonding curve
# - vTokensInBondingCurve: Virtual tokens in bonding curve
# - signature: Transaction signature
if not all([name, symbol, mint_str, bonding_curve_str, creator_str]):
logger.warning("Missing required fields in PumpPortal token data")
return None
# Convert string addresses to Pubkey objects
mint = Pubkey.from_string(mint_str)
bonding_curve = Pubkey.from_string(bonding_curve_str)
user = Pubkey.from_string(creator_str)
# For PumpPortal, we assume the creator is the same as the user
# since PumpPortal doesn't distinguish between them
creator = user
# Calculate derived addresses
associated_bonding_curve = self._find_associated_bonding_curve(
mint, bonding_curve
)
creator_vault = self._find_creator_vault(creator)
return TokenInfo(
name=name,
symbol=symbol,
uri=uri,
mint=mint,
platform=Platform.PUMP_FUN,
bonding_curve=bonding_curve,
associated_bonding_curve=associated_bonding_curve,
user=user,
creator=creator,
creator_vault=creator_vault,
)
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
@@ -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")