using System; using MongoDB.Driver; using PolyTraderSharp.Extensions; using System.Linq; using System.Net.WebSockets; using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Hosting; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services { public class AlchemyWebsocketService : BackgroundService { private const string CtfContractAddress = "0x4D97DCd97eC945f40cF65F87097ACe5EA0476045"; private const string TransferSingleTopic = "0xc3d58168c5ae7397731d063d5bbf3d657854427343f4c083240f7aacaa2d0f62"; private const string TransferBatchTopic = "0x4a39dc06d4c0dbc64b70af90fd698a233a518aa5d07e595d983b8c0526c8f7ce"; private readonly TradingState _state; private readonly CopyTradingState _copyState; private readonly ServerSettings _settings; private readonly TraderMonitorService _traderMonitor; private readonly TerminalLogger _logger; public AlchemyWebsocketService( TradingState state, CopyTradingState copyState, ServerSettings settings, TraderMonitorService traderMonitor, TerminalLogger logger) { _state = state; _copyState = copyState; _settings = settings; _traderMonitor = traderMonitor; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { if (!_settings.EnableBlockchainListener || string.IsNullOrEmpty(_settings.PolygonRpcUrl)) { _logger.Info("Blockchain Listener is disabled in settings. Using raw polling."); _state.IsAlchemyHealthy = false; return; } _logger.Info("Alchemy WSS Service starting up..."); while (!stoppingToken.IsCancellationRequested) { if (_state.GlobalTradingPaused || (_state.LiveTradingMode == TradingMode.Inactive && _state.DemoTradingMode == TradingMode.Inactive)) { _state.IsAlchemyHealthy = false; await Task.Delay(5000, stoppingToken); continue; } try { await ConnectAndListenAsync(stoppingToken); } catch (WebSocketException ex) { // Usually indicates a connection drop or 429 _logger.Warning($"Alchemy WSS drop: {ex.Message}. Falling back to API polling for 5 minutes."); _state.IsAlchemyHealthy = false; await Task.Delay(TimeSpan.FromMinutes(5), stoppingToken); } catch (Exception ex) { if (ex.Message.Contains("429") || ex.Message.Contains("Too Many Requests")) { _logger.Error($"Alchemy HTTP 429 Limit reached. Suspending WSS for 5 minutes."); _state.IsAlchemyHealthy = false; await Task.Delay(TimeSpan.FromMinutes(5), stoppingToken); } else { _logger.Error($"Alchemy WSS Error: {ex.Message}. Retrying in 10s..."); _state.IsAlchemyHealthy = false; await Task.Delay(10000, stoppingToken); } } } } private async Task ConnectAndListenAsync(CancellationToken stoppingToken) { using var ws = new ClientWebSocket(); var wssUrl = _settings.PolygonRpcUrl.Replace("https://", "wss://").Replace("http://", "ws://"); _logger.Info($"Connecting to Alchemy WebSocket: {wssUrl.Substring(0, Math.Min(35, wssUrl.Length))}..."); await ws.ConnectAsync(new Uri(wssUrl), stoppingToken); _state.IsAlchemyHealthy = true; _logger.Info("✅ Alchemy WSS Connected. Dynamic API-Throttling activated."); var activeTraders = _copyState.Traders.Values.Where(t => t.IsActive).ToList(); var activeStateHash = string.Join(",", activeTraders.OrderBy(t => t.Id).Select(t => t.WalletAddress.ToLowerInvariant())); if (activeTraders.Count > 0) { var paddedAddresses = activeTraders.Select(t => PadAddress(t.WalletAddress)).ToList(); int batchSize = 3; // Alchemy limits Topic arrays to a max of 3/4 entries int reqId = 1; for (int i = 0; i < paddedAddresses.Count; i += batchSize) { var chunk = paddedAddresses.Skip(i).Take(batchSize).ToList(); var addrJson = JsonSerializer.Serialize(chunk); // Buys: Master Trader is the receiver (Topic 3) var subscribeBuysStr = $@"{{ ""jsonrpc"": ""2.0"", ""id"": {reqId++}, ""method"": ""eth_subscribe"", ""params"": [ ""logs"", {{ ""address"": ""{CtfContractAddress}"", ""topics"": [ [""{TransferSingleTopic}"", ""{TransferBatchTopic}""], null, null, {addrJson} ] }} ] }}"; // Sells: Master Trader is the sender (Topic 2) var subscribeSellsStr = $@"{{ ""jsonrpc"": ""2.0"", ""id"": {reqId++}, ""method"": ""eth_subscribe"", ""params"": [ ""logs"", {{ ""address"": ""{CtfContractAddress}"", ""topics"": [ [""{TransferSingleTopic}"", ""{TransferBatchTopic}""], null, {addrJson} ] }} ] }}"; await ws.SendAsync(new ArraySegment(Encoding.UTF8.GetBytes(subscribeBuysStr)), WebSocketMessageType.Text, true, stoppingToken); await ws.SendAsync(new ArraySegment(Encoding.UTF8.GetBytes(subscribeSellsStr)), WebSocketMessageType.Text, true, stoppingToken); } } else { _logger.Info("Keine aktiven Master-Trader. Socket läuft im Standby..."); } var buffer = new byte[1024 * 64]; var loopCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); var monitorTask = Task.Run(async () => { while (!loopCts.IsCancellationRequested) { await Task.Delay(5000, loopCts.Token); var currentTraders = _copyState.Traders.Values.Where(t => t.IsActive).ToList(); var currentHash = string.Join(",", currentTraders.OrderBy(t => t.Id).Select(t => t.WalletAddress.ToLowerInvariant())); if (currentHash != activeStateHash) { _logger.Info("🔄 Master-Trader Liste hat sich geändert. Starte Alchemy WSS mit neuen Filtern neu..."); loopCts.Cancel(); break; } } }); try { while (ws.State == WebSocketState.Open && !loopCts.IsCancellationRequested) { var result = await ws.ReceiveAsync(new ArraySegment(buffer), loopCts.Token); if (result.MessageType == WebSocketMessageType.Close) { break; } var message = Encoding.UTF8.GetString(buffer, 0, result.Count); try { ProcessMessage(message); } catch (Exception ex) { _logger.Error($"Error parsing WSS msg: {ex.Message}"); } } } catch (OperationCanceledException) { // Expected when reconnecting due to trader list change } finally { if (!loopCts.IsCancellationRequested) loopCts.Cancel(); } } private void ProcessMessage(string jsonStr) { using var doc = JsonDocument.Parse(jsonStr); var root = doc.RootElement; if (!root.TryGetProperty("params", out var paramsEl)) return; if (!paramsEl.TryGetProperty("result", out var resultEl)) return; if (!resultEl.TryGetProperty("topics", out var topicsEl) || topicsEl.GetArrayLength() < 4) return; var topics = topicsEl.EnumerateArray().Select(t => t.GetString()).ToList(); var fromTopic = topics[2]?.ToLowerInvariant(); var toTopic = topics[3]?.ToLowerInvariant(); if (fromTopic == null || toTopic == null) return; var activeTraders = _copyState.Traders.Values.Where(t => t.IsActive).ToList(); string? triggeredAddress = null; foreach (var trader in activeTraders) { var padded = PadAddress(trader.WalletAddress); if (fromTopic == padded || toTopic == padded) { triggeredAddress = trader.WalletAddress; break; } } if (!string.IsNullOrEmpty(triggeredAddress)) { string txHash = resultEl.TryGetProperty("transactionHash", out var th) ? th.GetString() ?? "unknown" : "unknown"; _traderMonitor.TriggerFastBlockchainPoll(txHash, _settings.PolygonRpcUrl, triggeredAddress); } } private string PadAddress(string address) { string stripped = address.Replace("0x", "", StringComparison.OrdinalIgnoreCase).ToLowerInvariant(); return "0x" + stripped.PadLeft(64, '0'); } } }