Files
pumpfun-bonkfun-bot/src/trading/trader.py
T

295 lines
10 KiB
Python
Raw Normal View History

2025-03-05 16:44:55 +00:00
"""
Main trading coordinator for pump.fun tokens.
2025-03-06 16:14:17 +00:00
Refactored PumpTrader to only process fresh tokens from WebSocket.
2025-03-05 16:44:55 +00:00
"""
import asyncio
import json
import os
from datetime import datetime
2025-03-06 13:52:40 +00:00
import config
from core.priority_fee.manager import PriorityFeeManager
2025-03-05 16:44:55 +00:00
from src.core.client import SolanaClient
from src.core.curve import BondingCurveManager
from src.core.pubkeys import PumpAddresses
from src.core.wallet import Wallet
from src.monitoring.listener import PumpTokenListener
from src.trading.base import TokenInfo, TradeResult
from src.trading.buyer import TokenBuyer
from src.trading.seller import TokenSeller
from src.utils.logger import get_logger
logger = get_logger(__name__)
class PumpTrader:
2025-03-06 16:14:17 +00:00
"""Coordinates trading operations for pump.fun tokens with focus on freshness."""
2025-03-05 16:44:55 +00:00
def __init__(
self,
rpc_endpoint: str,
wss_endpoint: str,
private_key: str,
buy_amount: float,
buy_slippage: float,
sell_slippage: float,
max_retries: int = 5,
):
"""Initialize the pump trader.
Args:
rpc_endpoint: RPC endpoint URL
wss_endpoint: WebSocket endpoint URL
private_key: Wallet private key
buy_amount: Amount of SOL to spend on buys
buy_slippage: Slippage tolerance for buys
sell_slippage: Slippage tolerance for sells
max_retries: Maximum number of retry attempts
"""
self.solana_client = SolanaClient(rpc_endpoint)
self.wallet = Wallet(private_key)
self.curve_manager = BondingCurveManager(self.solana_client)
self.priority_fee_manager = PriorityFeeManager(
client=self.solana_client,
enable_dynamic_fee=config.ENABLE_DYNAMIC_PRIORITY_FEE,
enable_fixed_fee=config.ENABLE_FIXED_PRIORITY_FEE,
fixed_fee=config.FIXED_PRIORITY_FEE,
extra_fee=config.EXTRA_PRIORITY_FEE,
hard_cap=config.HARD_CAP_PRIOR_FEE,
)
2025-03-05 16:44:55 +00:00
self.buyer = TokenBuyer(
self.solana_client,
self.wallet,
self.curve_manager,
self.priority_fee_manager,
2025-03-05 16:44:55 +00:00
buy_amount,
buy_slippage,
max_retries,
)
self.seller = TokenSeller(
self.solana_client,
self.wallet,
self.curve_manager,
self.priority_fee_manager,
2025-03-05 16:44:55 +00:00
sell_slippage,
max_retries,
)
self.token_listener = PumpTokenListener(wss_endpoint, PumpAddresses.PROGRAM)
self.buy_amount = buy_amount
self.buy_slippage = buy_slippage
self.sell_slippage = sell_slippage
self.max_retries = max_retries
2025-03-06 22:05:51 +00:00
self.max_token_age = config.MAX_TOKEN_AGE
2025-03-06 16:14:17 +00:00
# Token processing state
self.token_queue = asyncio.Queue()
self.processing = False
2025-03-06 21:52:43 +00:00
self.processed_tokens: set[str] = set()
self.token_timestamps: dict[str, float] = {}
2025-03-05 16:44:55 +00:00
async def start(
self,
match_string: str | None = None,
bro_address: str | None = None,
marry_mode: bool = False,
yolo_mode: bool = False,
) -> None:
"""Start the trading bot.
Args:
match_string: Optional string to match in token name/symbol
bro_address: Optional creator address to filter by
marry_mode: If True, only buy tokens and skip selling
yolo_mode: If True, trade continuously
"""
logger.info("Starting pump.fun trader")
logger.info(f"Match filter: {match_string if match_string else 'None'}")
logger.info(f"Creator filter: {bro_address if bro_address else 'None'}")
logger.info(f"Marry mode: {marry_mode}")
logger.info(f"YOLO mode: {yolo_mode}")
2025-03-06 16:14:17 +00:00
logger.info(f"Max token age: {self.max_token_age} seconds")
# Start processor task
processor_task = asyncio.create_task(
self._process_token_queue(marry_mode, yolo_mode)
)
2025-03-05 16:44:55 +00:00
try:
await self.token_listener.listen_for_tokens(
2025-03-06 16:14:17 +00:00
lambda token: self._queue_token(token),
2025-03-05 16:44:55 +00:00
match_string,
bro_address,
)
except Exception as e:
logger.error(f"Trading stopped due to error: {str(e)}")
2025-03-06 16:14:17 +00:00
processor_task.cancel()
2025-03-05 16:44:55 +00:00
await self.solana_client.close()
2025-03-06 16:14:17 +00:00
async def _queue_token(self, token_info: TokenInfo) -> None:
"""Queue a token for processing if not already processed."""
token_key = str(token_info.mint)
if token_key in self.processed_tokens:
logger.debug(f"Token {token_info.symbol} already processed. Skipping...")
return
# Record timestamp when token was discovered
self.token_timestamps[token_key] = asyncio.get_event_loop().time()
await self.token_queue.put(token_info)
2025-03-06 22:05:51 +00:00
# logger.info(f"Queued new token: {token_info.symbol} ({token_info.mint})")
2025-03-06 16:14:17 +00:00
async def _process_token_queue(self, marry_mode: bool, yolo_mode: bool) -> None:
"""Continuously process tokens from the queue, only if they're fresh."""
while True:
token_info = await self.token_queue.get()
token_key = str(token_info.mint)
2025-03-06 22:05:51 +00:00
# Check if token is still "fresh"
2025-03-06 16:14:17 +00:00
current_time = asyncio.get_event_loop().time()
token_age = current_time - self.token_timestamps.get(
token_key, current_time
)
if token_age > self.max_token_age:
logger.info(
f"Skipping token {token_info.symbol} - too old ({token_age:.1f}s > {self.max_token_age}s)"
)
self.token_queue.task_done()
continue
self.processed_tokens.add(token_key)
logger.info(
f"Processing fresh token: {token_info.symbol} (age: {token_age:.1f}s)"
)
await self._handle_token(token_info, marry_mode, yolo_mode)
self.token_queue.task_done()
async def _handle_token(
2025-03-05 16:44:55 +00:00
self, token_info: TokenInfo, marry_mode: bool, yolo_mode: bool
) -> None:
"""Handle a new token creation event.
Args:
token_info: Token information
marry_mode: If True, only buy tokens and skip selling
yolo_mode: If True, continue trading after this token
"""
try:
await self._save_token_info(token_info)
2025-03-06 13:52:40 +00:00
logger.info(
f"Waiting for {config.WAIT_TIME_AFTER_CREATION} seconds for the bonding curve to stabilize..."
)
await asyncio.sleep(config.WAIT_TIME_AFTER_CREATION)
2025-03-05 16:44:55 +00:00
logger.info(
f"Buying {self.buy_amount:.6f} SOL worth of {token_info.symbol}..."
)
2025-03-06 13:52:40 +00:00
buy_result: TradeResult = await self.buyer.execute(token_info)
2025-03-05 16:44:55 +00:00
if buy_result.success:
logger.info(f"Successfully bought {token_info.symbol}")
2025-03-06 13:52:40 +00:00
self._log_trade(
"buy",
token_info,
buy_result.price, # type: ignore
buy_result.amount, # type: ignore
buy_result.tx_signature,
)
2025-03-05 16:44:55 +00:00
else:
logger.error(
f"Failed to buy {token_info.symbol}: {buy_result.error_message}"
)
# Sell token if not in marry mode
if not marry_mode and buy_result.success:
2025-03-06 13:52:40 +00:00
logger.info(
f"Waiting for {config.WAIT_TIME_AFTER_BUY} seconds before selling..."
)
await asyncio.sleep(config.WAIT_TIME_AFTER_BUY)
2025-03-05 16:44:55 +00:00
logger.info(f"Selling {token_info.symbol}...")
2025-03-06 13:52:40 +00:00
sell_result: TradeResult = await self.seller.execute(token_info)
2025-03-05 16:44:55 +00:00
if sell_result.success:
logger.info(f"Successfully sold {token_info.symbol}")
self._log_trade(
2025-03-06 13:52:40 +00:00
"sell",
token_info,
sell_result.price, # type: ignore
sell_result.amount, # type: ignore
sell_result.tx_signature,
2025-03-05 16:44:55 +00:00
)
else:
logger.error(
f"Failed to sell {token_info.symbol}: {sell_result.error_message}"
)
elif marry_mode:
logger.info("Marry mode enabled. Skipping sell operation.")
2025-03-06 16:14:17 +00:00
# Wait before looking for the next token
if yolo_mode:
logger.info(
f"YOLO mode enabled. Waiting {config.WAIT_TIME_BEFORE_NEW_TOKEN} seconds before looking for next token..."
)
await asyncio.sleep(config.WAIT_TIME_BEFORE_NEW_TOKEN)
2025-03-05 16:44:55 +00:00
except Exception as e:
logger.error(f"Error handling token {token_info.symbol}: {str(e)}")
async def _save_token_info(self, token_info: TokenInfo) -> None:
"""Save token information to a file.
Args:
token_info: Token information
"""
os.makedirs("trades", exist_ok=True)
file_name = os.path.join("trades", f"{token_info.mint}.txt")
with open(file_name, "w") as file:
file.write(json.dumps(token_info.to_dict(), indent=2))
logger.info(f"Token information saved to {file_name}")
def _log_trade(
2025-03-06 13:52:40 +00:00
self,
action: str,
token_info: TokenInfo,
price: float,
amount: float,
tx_hash: str | None,
2025-03-05 16:44:55 +00:00
) -> None:
"""Log trade information.
Args:
action: Trade action (buy/sell)
token_info: Token information
price: Token price in SOL
2025-03-06 13:52:40 +00:00
amount: Trade amount in SOL
2025-03-05 16:44:55 +00:00
tx_hash: Transaction hash
"""
os.makedirs("trades", exist_ok=True)
log_entry = {
"timestamp": datetime.utcnow().isoformat(),
"action": action,
"token_address": str(token_info.mint),
"symbol": token_info.symbol,
"price": price,
2025-03-06 13:52:40 +00:00
"amount": amount,
2025-03-06 16:14:17 +00:00
"tx_hash": str(tx_hash),
2025-03-05 16:44:55 +00:00
}
with open("trades/trades.log", "a") as log_file:
log_file.write(json.dumps(log_entry) + "\n")