diff --git a/Program.cs b/Program.cs index 7f0df7d..8adbf43 100644 --- a/Program.cs +++ b/Program.cs @@ -9,6 +9,7 @@ using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Options; using PolyTrader.Core.Configuration; using PolyTrader.Core.DependencyInjection; +using PolyTrader.Core.Streaming; using PolyTraderSharp.Models; using PolyTraderSharp.Services; @@ -35,6 +36,7 @@ internal static class Program return client.GetDatabase(dbOptions.DatabaseName); }); services.AddCorePersistence(); + services.AddSingleton(); services.AddSingleton((IServiceProvider sp) => ServerSettings.Load("server_settings.xml")); services.AddSingleton(); services.AddSingleton(); diff --git a/services/AlchemyWebsocketService.cs b/services/AlchemyWebsocketService.cs index e6387eb..25114e2 100644 --- a/services/AlchemyWebsocketService.cs +++ b/services/AlchemyWebsocketService.cs @@ -1,13 +1,11 @@ using System; -using MongoDB.Driver; -using PolyTraderSharp.Extensions; +using System.Collections.Generic; 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 PolyTrader.Core.Streaming; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services @@ -23,19 +21,22 @@ namespace PolyTraderSharp.Services private readonly ServerSettings _settings; private readonly TraderMonitorService _traderMonitor; private readonly TerminalLogger _logger; + private readonly IBlockchainWssClientFactory _wssFactory; public AlchemyWebsocketService( TradingState state, CopyTradingState copyState, ServerSettings settings, TraderMonitorService traderMonitor, - TerminalLogger logger) + TerminalLogger logger, + IBlockchainWssClientFactory wssFactory) { _state = state; _copyState = copyState; _settings = settings; _traderMonitor = traderMonitor; _logger = logger; + _wssFactory = wssFactory; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) @@ -51,7 +52,7 @@ namespace PolyTraderSharp.Services while (!stoppingToken.IsCancellationRequested) { - if (_state.GlobalTradingPaused || + if (_state.GlobalTradingPaused || (_state.LiveTradingMode == TradingMode.Inactive && _state.DemoTradingMode == TradingMode.Inactive)) { _state.IsAlchemyHealthy = false; @@ -88,86 +89,39 @@ namespace PolyTraderSharp.Services } } + /// + /// Baut EINE WSS-Session über den Core-Client auf. Consumer-Verantwortung: Filter aus den + /// getrackten Master-Tradern bauen, auf Trader-Listen-Änderungen re-subscriben, Health spiegeln. + /// 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 + var subscription = BuildSubscription(activeTraders); + if (subscription.Filters.Count == 0) { _logger.Info("Keine aktiven Master-Trader. Socket läuft im Standby..."); + _state.IsAlchemyHealthy = false; + await Task.Delay(5000, stoppingToken); + return; } - var buffer = new byte[1024 * 64]; + var client = _wssFactory.Create(); + _logger.Info("Connecting to Alchemy WebSocket (Core WSS client)..."); - var loopCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); - - var monitorTask = Task.Run(async () => + using var loopCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); + + // Überwacht Trader-Listen-Änderungen (Re-Subscribe) und spiegelt das Health-Flag. + var monitorTask = Task.Run(async () => { while (!loopCts.IsCancellationRequested) { - await Task.Delay(5000, loopCts.Token); + try { await Task.Delay(5000, loopCts.Token); } + catch (OperationCanceledException) { break; } + + _state.IsAlchemyHealthy = client.IsHealthy; + 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) @@ -181,49 +135,43 @@ namespace PolyTraderSharp.Services 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}"); - } - } + await client.ConnectAndListenAsync(_settings.PolygonRpcUrl, subscription, OnLog, loopCts.Token); } catch (OperationCanceledException) { - // Expected when reconnecting due to trader list change + // Erwartet beim Re-Subscribe (Trader-Liste geändert) oder beim Shutdown. } finally { + _state.IsAlchemyHealthy = false; if (!loopCts.IsCancellationRequested) loopCts.Cancel(); } } - private void ProcessMessage(string jsonStr) + private BlockchainWssSubscription BuildSubscription(List activeTraders) { - using var doc = JsonDocument.Parse(jsonStr); - var root = doc.RootElement; + var subscription = new BlockchainWssSubscription(); + var paddedAddresses = activeTraders.Select(t => PadAddress(t.WalletAddress)).ToList(); + var topic0 = new List { TransferSingleTopic, TransferBatchTopic }; - 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; + const int batchSize = 3; // Alchemy begrenzt Topic-Arrays auf max. 3-4 Einträge + for (int i = 0; i < paddedAddresses.Count; i += batchSize) + { + var chunk = paddedAddresses.Skip(i).Take(batchSize).ToList(); + // Buys: Master-Trader ist Empfänger (Topic 3) + subscription.Filters.Add(new LogSubscriptionFilter { Address = CtfContractAddress, Topic0 = topic0, Topic3 = chunk }); + // Sells: Master-Trader ist Sender (Topic 2) + subscription.Filters.Add(new LogSubscriptionFilter { Address = CtfContractAddress, Topic0 = topic0, Topic2 = chunk }); + } + return subscription; + } - var topics = topicsEl.EnumerateArray().Select(t => t.GetString()).ToList(); - var fromTopic = topics[2]?.ToLowerInvariant(); - var toTopic = topics[3]?.ToLowerInvariant(); + private void OnLog(BlockchainLogEvent evt) + { + if (evt.Topics.Count < 4) return; + var fromTopic = evt.Topics[2]?.ToLowerInvariant(); + var toTopic = evt.Topics[3]?.ToLowerInvariant(); if (fromTopic == null || toTopic == null) return; var activeTraders = _copyState.Traders.Values.Where(t => t.IsActive).ToList(); @@ -241,8 +189,7 @@ namespace PolyTraderSharp.Services if (!string.IsNullOrEmpty(triggeredAddress)) { - string txHash = resultEl.TryGetProperty("transactionHash", out var th) ? th.GetString() ?? "unknown" : "unknown"; - _traderMonitor.TriggerFastBlockchainPoll(txHash, _settings.PolygonRpcUrl, triggeredAddress); + _traderMonitor.TriggerFastBlockchainPoll(evt.TransactionHash, _settings.PolygonRpcUrl, triggeredAddress); } } diff --git a/src/PolyTrader.Core/Streaming/AlchemyWssClient.cs b/src/PolyTrader.Core/Streaming/AlchemyWssClient.cs new file mode 100644 index 0000000..fa68451 --- /dev/null +++ b/src/PolyTrader.Core/Streaming/AlchemyWssClient.cs @@ -0,0 +1,103 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Net.WebSockets; +using System.Text; +using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; + +namespace PolyTrader.Core.Streaming +{ + /// + /// EVM/Alchemy-WSS-Verbindungsklasse (eth_subscribe "logs"). Domänen-agnostisch: + /// Contract/Topics/Adress-Filter kommen vom Aufrufer über die Subscription. + /// + public class AlchemyWssClient : IBlockchainWssClient + { + public bool IsHealthy { get; private set; } + + public async Task ConnectAndListenAsync(string rpcUrl, BlockchainWssSubscription subscription, + Action onLog, CancellationToken ct) + { + using var ws = new ClientWebSocket(); + var wssUrl = rpcUrl.Replace("https://", "wss://").Replace("http://", "ws://"); + + await ws.ConnectAsync(new Uri(wssUrl), ct); + IsHealthy = true; + + try + { + int reqId = 1; + foreach (var filter in subscription.Filters) + { + var payload = BuildSubscribePayload(reqId++, filter); + await ws.SendAsync(new ArraySegment(Encoding.UTF8.GetBytes(payload)), + WebSocketMessageType.Text, true, ct); + } + + var buffer = new byte[1024 * 64]; + while (ws.State == WebSocketState.Open && !ct.IsCancellationRequested) + { + var result = await ws.ReceiveAsync(new ArraySegment(buffer), ct); + if (result.MessageType == WebSocketMessageType.Close) break; + + var message = Encoding.UTF8.GetString(buffer, 0, result.Count); + try + { + var evt = TryParseLog(message); + if (evt != null) onLog(evt); + } + catch + { + // Resilient: eine fehlerhafte Nachricht oder ein Handler-Fehler + // darf die Session nicht beenden (Verhalten wie zuvor). + } + } + } + finally + { + IsHealthy = false; + } + } + + private static string BuildSubscribePayload(int id, LogSubscriptionFilter f) + { + // Positionsbasiertes Topics-Array; trailing-Wildcards (null) werden abgeschnitten, + // exakt wie im ursprünglichen Handbuild (Sells: 3 Positionen, Buys: 4). + var positions = new List + { + JsonSerializer.Serialize(f.Topic0), + f.Topic1 != null ? JsonSerializer.Serialize(f.Topic1) : "null", + f.Topic2 != null ? JsonSerializer.Serialize(f.Topic2) : "null", + f.Topic3 != null ? JsonSerializer.Serialize(f.Topic3) : "null" + }; + while (positions.Count > 1 && positions[^1] == "null") positions.RemoveAt(positions.Count - 1); + + string topicsJson = "[" + string.Join(",", positions) + "]"; + return $@"{{""jsonrpc"":""2.0"",""id"":{id},""method"":""eth_subscribe"",""params"":[""logs"",{{""address"":""{f.Address}"",""topics"":{topicsJson}}}]}}"; + } + + private static BlockchainLogEvent? TryParseLog(string jsonStr) + { + using var doc = JsonDocument.Parse(jsonStr); + var root = doc.RootElement; + + if (!root.TryGetProperty("params", out var paramsEl)) return null; + if (!paramsEl.TryGetProperty("result", out var resultEl)) return null; + if (!resultEl.TryGetProperty("topics", out var topicsEl)) return null; + + var topics = topicsEl.EnumerateArray().Select(t => t.GetString()).ToList(); + string txHash = resultEl.TryGetProperty("transactionHash", out var th) ? th.GetString() ?? "" : ""; + string address = resultEl.TryGetProperty("address", out var ad) ? ad.GetString() ?? "" : ""; + string data = resultEl.TryGetProperty("data", out var dt) ? dt.GetString() ?? "" : ""; + + return new BlockchainLogEvent { Topics = topics, TransactionHash = txHash, Address = address, Data = data }; + } + } + + public class AlchemyWssClientFactory : IBlockchainWssClientFactory + { + public IBlockchainWssClient Create() => new AlchemyWssClient(); + } +} diff --git a/src/PolyTrader.Core/Streaming/BlockchainWss.cs b/src/PolyTrader.Core/Streaming/BlockchainWss.cs new file mode 100644 index 0000000..c7f3398 --- /dev/null +++ b/src/PolyTrader.Core/Streaming/BlockchainWss.cs @@ -0,0 +1,68 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +namespace PolyTrader.Core.Streaming +{ + /// + /// Ein einzelner eth_subscribe("logs")-Filter. Positionsbasiert wie bei Alchemy: + /// Topic0 = Event-Signaturen (OR-Array), Topic1..3 = indizierte Parameter + /// (z.B. Adress-Filter). null an einer Position = Wildcard. + /// + public class LogSubscriptionFilter + { + public string Address { get; set; } = string.Empty; + public List Topic0 { get; set; } = new(); + public List? Topic1 { get; set; } + public List? Topic2 { get; set; } + public List? Topic3 { get; set; } + } + + /// + /// Beschreibt, was ein WSS-Client abonniert. Der Aufrufer (Modul) baut die Filter + /// aus seinen eigenen Daten (z.B. Wallets der getrackten Master-Trader) und batcht sie. + /// + public class BlockchainWssSubscription + { + public List Filters { get; } = new(); + } + + /// + /// Ein dekodiertes Log-Event aus dem Stream (domänen-agnostisch). + /// Die fachliche Auswertung (welche Adresse, welche Reaktion) macht der Aufrufer. + /// + public class BlockchainLogEvent + { + public IReadOnlyList Topics { get; init; } = Array.Empty(); + public string TransactionHash { get; init; } = string.Empty; + public string Address { get; init; } = string.Empty; + public string Data { get; init; } = string.Empty; + } + + /// + /// Wiederverwendbare WSS-Verbindungsfähigkeit (Core-Infrastruktur). Kapselt Connect, + /// eth_subscribe, Empfangs-Loop und Decode. Reconnect-Policy, Gating und der Filterinhalt + /// liegen beim Aufrufer, damit jedes Modul eine eigene, gefilterte Instanz mit eigenem + /// Key betreiben kann (kein geteilter, gefluteter Singleton-Stream). + /// + public interface IBlockchainWssClient + { + /// True während einer aktiven, verbundenen Session. + bool IsHealthy { get; } + + /// + /// Baut EINE Session auf: verbindet, sendet die Subscriptions, liest bis Close/Cancel/Fehler + /// und ruft je dekodiertem Log. Wirft bei Verbindungsfehlern — + /// die Reconnect-Entscheidung trifft der Aufrufer. + /// + Task ConnectAndListenAsync(string rpcUrl, BlockchainWssSubscription subscription, + Action onLog, CancellationToken ct); + } + + /// Erzeugt pro Aufrufer/Modul eine eigene WSS-Client-Instanz. + public interface IBlockchainWssClientFactory + { + IBlockchainWssClient Create(); + } +}