mirror of
https://github.com/chainstacklabs/pumpfun-bonkfun-bot.git
synced 2026-08-07 12:37:47 +00:00
docs(claude): add claude code rules
This commit is contained in:
@@ -13,7 +13,7 @@ class BaseTokenListener(ABC):
|
||||
|
||||
def __init__(self, platform: Platform | None = None):
|
||||
"""Initialize the listener with optional platform specification.
|
||||
|
||||
|
||||
Args:
|
||||
platform: Platform to monitor (if None, monitor all platforms)
|
||||
"""
|
||||
@@ -38,13 +38,13 @@ class BaseTokenListener(ABC):
|
||||
|
||||
def should_process_token(self, token_info: TokenInfo) -> bool:
|
||||
"""Check if a token should be processed based on platform filter.
|
||||
|
||||
|
||||
Args:
|
||||
token_info: Token information
|
||||
|
||||
|
||||
Returns:
|
||||
True if token should be processed
|
||||
"""
|
||||
if self.platform is None:
|
||||
return True # Process all platforms
|
||||
return token_info.platform == self.platform
|
||||
return token_info.platform == self.platform
|
||||
|
||||
@@ -48,7 +48,7 @@ class ListenerFactory:
|
||||
)
|
||||
|
||||
from monitoring.universal_geyser_listener import UniversalGeyserListener
|
||||
|
||||
|
||||
listener = UniversalGeyserListener(
|
||||
geyser_endpoint=geyser_endpoint,
|
||||
geyser_api_token=geyser_api_token,
|
||||
@@ -63,7 +63,7 @@ class ListenerFactory:
|
||||
raise ValueError("WebSocket endpoint is required for logs listener")
|
||||
|
||||
from monitoring.universal_logs_listener import UniversalLogsListener
|
||||
|
||||
|
||||
listener = UniversalLogsListener(
|
||||
wss_endpoint=wss_endpoint,
|
||||
platforms=platforms,
|
||||
@@ -76,7 +76,7 @@ class ListenerFactory:
|
||||
raise ValueError("WebSocket endpoint is required for blocks listener")
|
||||
|
||||
from monitoring.universal_block_listener import UniversalBlockListener
|
||||
|
||||
|
||||
listener = UniversalBlockListener(
|
||||
wss_endpoint=wss_endpoint,
|
||||
platforms=platforms,
|
||||
@@ -89,26 +89,36 @@ class ListenerFactory:
|
||||
from monitoring.universal_pumpportal_listener import (
|
||||
UniversalPumpPortalListener,
|
||||
)
|
||||
|
||||
|
||||
# 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]
|
||||
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")
|
||||
|
||||
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]
|
||||
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")
|
||||
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)]}")
|
||||
logger.info(
|
||||
f"Created Universal PumpPortal listener for platforms: {[p.value for p in (platforms or supported_pumpportal_platforms)]}"
|
||||
)
|
||||
return listener
|
||||
|
||||
else:
|
||||
@@ -146,8 +156,8 @@ class ListenerFactory:
|
||||
@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]
|
||||
return [Platform.PUMP_FUN, Platform.LETS_BONK]
|
||||
|
||||
@@ -33,25 +33,25 @@ class UniversalBlockListener(BaseTokenListener):
|
||||
super().__init__()
|
||||
self.wss_endpoint = wss_endpoint
|
||||
self.ping_interval = 20 # seconds
|
||||
|
||||
|
||||
# Import platform factory and get supported platforms
|
||||
from platforms import platform_factory
|
||||
|
||||
|
||||
if platforms is None:
|
||||
# Monitor all supported platforms
|
||||
self.platforms = platform_factory.get_supported_platforms()
|
||||
else:
|
||||
self.platforms = platforms
|
||||
|
||||
|
||||
# Get event parsers for all platforms
|
||||
self.platform_parsers = {}
|
||||
self.platform_program_ids = []
|
||||
|
||||
|
||||
for platform in self.platforms:
|
||||
try:
|
||||
# Create a simple dummy client that doesn't start blockhash updater
|
||||
from core.client import SolanaClient
|
||||
|
||||
|
||||
# Create a mock client class to avoid network operations
|
||||
class DummyClient(SolanaClient):
|
||||
def __init__(self):
|
||||
@@ -61,16 +61,18 @@ class UniversalBlockListener(BaseTokenListener):
|
||||
self._cached_blockhash = None
|
||||
self._blockhash_lock = None
|
||||
self._blockhash_updater_task = None
|
||||
|
||||
|
||||
dummy_client = DummyClient()
|
||||
|
||||
|
||||
implementations = get_platform_implementations(platform, dummy_client)
|
||||
parser = implementations.event_parser
|
||||
self.platform_parsers[platform] = parser
|
||||
self.platform_program_ids.append(str(parser.get_program_id()))
|
||||
|
||||
logger.info(f"Registered platform {platform.value} with program ID {parser.get_program_id()}")
|
||||
|
||||
|
||||
logger.info(
|
||||
f"Registered platform {platform.value} with program ID {parser.get_program_id()}"
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning(f"Could not register platform {platform.value}: {e}")
|
||||
|
||||
@@ -117,7 +119,10 @@ class UniversalBlockListener(BaseTokenListener):
|
||||
)
|
||||
continue
|
||||
|
||||
if creator_address and str(token_info.user) != creator_address:
|
||||
if (
|
||||
creator_address
|
||||
and str(token_info.user) != creator_address
|
||||
):
|
||||
logger.info(
|
||||
f"Token not created by {creator_address}. Skipping..."
|
||||
)
|
||||
@@ -220,10 +225,10 @@ class UniversalBlockListener(BaseTokenListener):
|
||||
|
||||
for platform, parser in self.platform_parsers.items():
|
||||
# Check if the parser has a block parsing method
|
||||
if hasattr(parser, 'parse_token_creation_from_block'):
|
||||
token_info = parser.parse_token_creation_from_block({
|
||||
"transactions": [tx]
|
||||
})
|
||||
if hasattr(parser, "parse_token_creation_from_block"):
|
||||
token_info = parser.parse_token_creation_from_block(
|
||||
{"transactions": [tx]}
|
||||
)
|
||||
if token_info:
|
||||
return token_info
|
||||
|
||||
@@ -237,4 +242,4 @@ class UniversalBlockListener(BaseTokenListener):
|
||||
except Exception:
|
||||
logger.exception("Error processing WebSocket message")
|
||||
|
||||
return None
|
||||
return None
|
||||
|
||||
@@ -30,7 +30,7 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
super().__init__()
|
||||
self.geyser_endpoint = geyser_endpoint
|
||||
self.geyser_api_token = geyser_api_token
|
||||
|
||||
|
||||
valid_auth_types = {"x-token", "basic"}
|
||||
self.auth_type: str = (geyser_auth_type or "x-token").lower()
|
||||
if self.auth_type not in valid_auth_types:
|
||||
@@ -38,21 +38,21 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
f"Unsupported auth_type={self.auth_type!r}. "
|
||||
f"Expected one of {valid_auth_types}"
|
||||
)
|
||||
|
||||
|
||||
if platforms is None:
|
||||
self.platforms = platform_factory.get_supported_platforms()
|
||||
else:
|
||||
self.platforms = platforms
|
||||
|
||||
|
||||
# Get event parsers for all platforms
|
||||
self.platform_parsers = {}
|
||||
self.platform_program_ids = set()
|
||||
|
||||
|
||||
for platform in self.platforms:
|
||||
try:
|
||||
# Create a simple dummy client that doesn't start blockhash updater
|
||||
from core.client import SolanaClient
|
||||
|
||||
|
||||
# Create a mock client class to avoid network operations
|
||||
class DummyClient(SolanaClient):
|
||||
def __init__(self):
|
||||
@@ -62,22 +62,26 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
self._cached_blockhash = None
|
||||
self._blockhash_lock = None
|
||||
self._blockhash_updater_task = None
|
||||
|
||||
|
||||
dummy_client = DummyClient()
|
||||
|
||||
implementations = platform_factory.create_for_platform(platform, dummy_client)
|
||||
|
||||
implementations = platform_factory.create_for_platform(
|
||||
platform, dummy_client
|
||||
)
|
||||
parser = implementations.event_parser
|
||||
self.platform_parsers[platform] = parser
|
||||
self.platform_program_ids.add(parser.get_program_id())
|
||||
|
||||
logger.info(f"Registered platform {platform.value} with program ID {parser.get_program_id()}")
|
||||
|
||||
|
||||
logger.info(
|
||||
f"Registered platform {platform.value} with program ID {parser.get_program_id()}"
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning(f"Could not register platform {platform.value}: {e}")
|
||||
|
||||
async def _create_geyser_connection(self):
|
||||
"""Establish a secure connection to the Geyser endpoint."""
|
||||
|
||||
|
||||
if self.auth_type == "x-token":
|
||||
auth = grpc.metadata_call_credentials(
|
||||
lambda _, callback: callback(
|
||||
@@ -92,20 +96,20 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
)
|
||||
creds = grpc.composite_channel_credentials(grpc.ssl_channel_credentials(), auth)
|
||||
channel = grpc.aio.secure_channel(self.geyser_endpoint, creds)
|
||||
|
||||
|
||||
return geyser_pb2_grpc.GeyserStub(channel), channel
|
||||
|
||||
def _create_subscription_request(self):
|
||||
"""Create a subscription request for all monitored platforms."""
|
||||
|
||||
|
||||
request = geyser_pb2.SubscribeRequest()
|
||||
|
||||
|
||||
# Add all platform program IDs to the filter
|
||||
for program_id in self.platform_program_ids:
|
||||
filter_name = f"platform_filter_{program_id}"
|
||||
request.transactions[filter_name].account_include.append(str(program_id))
|
||||
request.transactions[filter_name].failed = False
|
||||
|
||||
|
||||
request.commitment = geyser_pb2.CommitmentLevel.PROCESSED
|
||||
return request
|
||||
|
||||
@@ -126,8 +130,12 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
request = self._create_subscription_request()
|
||||
|
||||
logger.info(f"Connected to Geyser endpoint: {self.geyser_endpoint}")
|
||||
logger.info(f"Monitoring platforms: {[p.value for p in self.platforms]}")
|
||||
logger.info(f"Monitoring program IDs: {[str(pid) for pid in self.platform_program_ids]}")
|
||||
logger.info(
|
||||
f"Monitoring platforms: {[p.value for p in self.platforms]}"
|
||||
)
|
||||
logger.info(
|
||||
f"Monitoring program IDs: {[str(pid) for pid in self.platform_program_ids]}"
|
||||
)
|
||||
|
||||
try:
|
||||
async for update in stub.Subscribe(iter([request])):
|
||||
@@ -192,7 +200,7 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
continue
|
||||
|
||||
program_id = Pubkey.from_bytes(msg.account_keys[program_idx])
|
||||
|
||||
|
||||
# Find the matching platform parser
|
||||
for platform, parser in self.platform_parsers.items():
|
||||
if program_id == parser.get_program_id():
|
||||
@@ -207,4 +215,4 @@ class UniversalGeyserListener(BaseTokenListener):
|
||||
|
||||
except Exception:
|
||||
logger.exception("Error processing Geyser update")
|
||||
return None
|
||||
return None
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""
|
||||
Universal logs listener that works with any platform through the interface system.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from collections.abc import Awaitable, Callable
|
||||
@@ -31,25 +32,25 @@ class UniversalLogsListener(BaseTokenListener):
|
||||
super().__init__()
|
||||
self.wss_endpoint = wss_endpoint
|
||||
self.ping_interval = 20 # seconds
|
||||
|
||||
|
||||
# Import platform factory and get supported platforms
|
||||
from platforms import platform_factory
|
||||
|
||||
|
||||
if platforms is None:
|
||||
# Monitor all supported platforms
|
||||
self.platforms = platform_factory.get_supported_platforms()
|
||||
else:
|
||||
self.platforms = platforms
|
||||
|
||||
|
||||
# Get event parsers for all platforms
|
||||
self.platform_parsers = {}
|
||||
self.platform_program_ids = []
|
||||
|
||||
|
||||
for platform in self.platforms:
|
||||
try:
|
||||
# Create a simple dummy client that doesn't start blockhash updater
|
||||
from core.client import SolanaClient
|
||||
|
||||
|
||||
# Create a mock client class to avoid network operations
|
||||
class DummyClient(SolanaClient):
|
||||
def __init__(self):
|
||||
@@ -59,16 +60,20 @@ class UniversalLogsListener(BaseTokenListener):
|
||||
self._cached_blockhash = None
|
||||
self._blockhash_lock = None
|
||||
self._blockhash_updater_task = None
|
||||
|
||||
|
||||
dummy_client = DummyClient()
|
||||
|
||||
implementations = platform_factory.create_for_platform(platform, dummy_client)
|
||||
|
||||
implementations = platform_factory.create_for_platform(
|
||||
platform, dummy_client
|
||||
)
|
||||
parser = implementations.event_parser
|
||||
self.platform_parsers[platform] = parser
|
||||
self.platform_program_ids.append(str(parser.get_program_id()))
|
||||
|
||||
logger.info(f"Registered platform {platform.value} with program ID {parser.get_program_id()}")
|
||||
|
||||
|
||||
logger.info(
|
||||
f"Registered platform {platform.value} with program ID {parser.get_program_id()}"
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning(f"Could not register platform {platform.value}: {e}")
|
||||
|
||||
@@ -115,7 +120,10 @@ class UniversalLogsListener(BaseTokenListener):
|
||||
)
|
||||
continue
|
||||
|
||||
if creator_address and str(token_info.user) != creator_address:
|
||||
if (
|
||||
creator_address
|
||||
and str(token_info.user) != creator_address
|
||||
):
|
||||
logger.info(
|
||||
f"Token not created by {creator_address}. Skipping..."
|
||||
)
|
||||
@@ -159,7 +167,9 @@ class UniversalLogsListener(BaseTokenListener):
|
||||
response = await websocket.recv()
|
||||
response_data = json.loads(response)
|
||||
if "result" in response_data:
|
||||
logger.info(f"Subscription confirmed with ID: {response_data['result']}")
|
||||
logger.info(
|
||||
f"Subscription confirmed with ID: {response_data['result']}"
|
||||
)
|
||||
else:
|
||||
logger.warning(f"Unexpected subscription response: {response}")
|
||||
|
||||
@@ -209,4 +219,4 @@ class UniversalLogsListener(BaseTokenListener):
|
||||
except Exception:
|
||||
logger.exception("Error processing WebSocket message")
|
||||
|
||||
return None
|
||||
return None
|
||||
|
||||
@@ -32,23 +32,23 @@ class UniversalPumpPortalListener(BaseTokenListener):
|
||||
super().__init__()
|
||||
self.pumpportal_url = pumpportal_url
|
||||
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:
|
||||
@@ -56,8 +56,10 @@ class UniversalPumpPortalListener(BaseTokenListener):
|
||||
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"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(
|
||||
@@ -100,8 +102,14 @@ class UniversalPumpPortalListener(BaseTokenListener):
|
||||
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 ""
|
||||
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..."
|
||||
@@ -197,7 +205,9 @@ class UniversalPumpPortalListener(BaseTokenListener):
|
||||
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")
|
||||
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}")
|
||||
@@ -213,4 +223,4 @@ class UniversalPumpPortalListener(BaseTokenListener):
|
||||
except Exception:
|
||||
logger.exception("Error processing PumpPortal WebSocket message")
|
||||
|
||||
return None
|
||||
return None
|
||||
|
||||
Reference in New Issue
Block a user