using System; using System.Collections.Concurrent; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; using Microsoft.Extensions.Hosting; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services { public class TraderMonitorService : BackgroundService { private readonly TradingState _state; private readonly PolymarketApiService _api; private readonly ChannelWriter _signalWriter; private readonly ChannelWriter _closedTradeWriter; private readonly TerminalLogger _logger; // Prevents duplicates. Fast O(1) lookup cache to prevent DB spam. private readonly ConcurrentDictionary _processedTxHashes = new(); private DateTime _lastHashCleanup = DateTime.UtcNow; private readonly ConcurrentDictionary _processedClosures = new(); private readonly ConcurrentDictionary _lastPolled = new(); private DateTime _lastLivePoll = DateTime.MinValue; public TraderMonitorService( TradingState state, PolymarketApiService api, ChannelWriter signalWriter, ChannelWriter closedTradeWriter, TerminalLogger logger) { _state = state; _api = api; _signalWriter = signalWriter; _closedTradeWriter = closedTradeWriter; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.Info("TraderMonitorService started background API priority polling..."); while (!stoppingToken.IsCancellationRequested) { try { await PollActiveTradersAsync(stoppingToken); // Live Accounts open positions sync (Runs every 30s instead of slamming API constantly) if ((DateTime.UtcNow - _lastLivePoll).TotalSeconds > 30) { await PollLiveAccountsAsync(stoppingToken); await PollDemoExpirationsAsync(stoppingToken); _lastLivePoll = DateTime.UtcNow; } } catch (Exception ex) { _logger.Error($"TraderMonitor polling error: {ex.Message}"); } // Global Engine Tick (dynamic queue evaluation) await Task.Delay(1000, stoppingToken); } } private async Task PollActiveTradersAsync(CancellationToken ct) { // Only process ACTIVE trader copies if not paused/inactive if (_state.GlobalTradingPaused || (_state.DemoTradingMode == TradingMode.Inactive && _state.LiveTradingMode == TradingMode.Inactive)) { return; } var activeTraders = _state.Traders.Values.Where(t => t.IsActive).ToList(); if (activeTraders.Count == 0) return; var now = DateTime.UtcNow; var toPoll = new List(); bool isWssHealthy = _state.IsAlchemyHealthy; // Calculate Dynamic Priorities // Data API rate limit: 1000 req/10s (general). // Worst case: 30 traders × high prio (3s) = ~100 req/10s = 10% capacity. // With medium prio at 10s and batches of 10: well within limits. foreach (var trader in activeTraders) { if (!_lastPolled.TryGetValue(trader.WalletAddress, out var lastPoll)) lastPoll = DateTime.MinValue; double secondsSinceLastPoll = (now - lastPoll).TotalSeconds; int requiredInterval = 10; // Medium Prio Default (Data API: 1000/10s headroom) if (isWssHealthy) { // If WSS is healthy, fall back to safety-net polling requiredInterval = 60; // 1 minute (was 2 min) } else { if (trader.TotalTrades > 20 || trader.Winrate30t >= 60.0) requiredInterval = 3; // High Prio (unchanged — already fast) else if (trader.TotalTrades < 5) requiredInterval = 30; // Low Prio (was 120s) } if (secondsSinceLastPoll >= requiredInterval) { toPoll.Add(trader); } } if (toPoll.Count == 0) return; // Batch Execution (Max 10 Concurrent Requests to respect API limits) int batchSize = 10; for (int i = 0; i < toPoll.Count; i += batchSize) { if (ct.IsCancellationRequested) break; var batch = toPoll.Skip(i).Take(batchSize); var tasks = batch.Select(async trader => { _lastPolled[trader.WalletAddress] = DateTime.UtcNow; System.Diagnostics.Stopwatch? sw = null; if (_state.DebugPollingLog) sw = System.Diagnostics.Stopwatch.StartNew(); var activity = await _api.GetTraderActivityAsync(trader.WalletAddress, limit: 50); if (_state.DebugPollingLog && sw != null) { sw.Stop(); _logger.Debug($"[API-Profiler] Activity-Request für Trader {trader.DisplayName} dauerte {sw.ElapsedMilliseconds} ms."); } foreach (var act in activity) { ProcessActivityItem(act, trader); } }); await Task.WhenAll(tasks); await Task.Delay(200, ct); // Tiny 200ms breath between batches } // Cleanup old hashes periodically (keep for 24 hours to prevent ANY duplicates) if ((DateTime.UtcNow - _lastHashCleanup).TotalHours > 1) { var cutoff = DateTime.UtcNow.AddHours(-24); var expired = _processedTxHashes.Where(x => x.Value < cutoff).Select(x => x.Key).ToList(); foreach (var k in expired) _processedTxHashes.TryRemove(k, out _); _lastHashCleanup = DateTime.UtcNow; } } /// /// Triggered instantly by the AlchemyWebsocketService when an EVM TransferSingle is detected. /// public void TriggerManualPoll(string walletAddress) { var trader = _state.Traders.Values.FirstOrDefault(t => t.WalletAddress.Equals(walletAddress, StringComparison.OrdinalIgnoreCase)); if (trader != null && trader.IsActive) { // Force an immediate poll on the next tick by artificially advancing the last poll date _lastPolled[trader.WalletAddress] = DateTime.MinValue; } } private async Task PollDemoExpirationsAsync(CancellationToken ct) { var demoAccounts = _state.Accounts.Values.Where(a => a.IsDemo && a.IsActive).ToList(); if (demoAccounts.Count == 0) return; foreach (var acc in demoAccounts) { if (ct.IsCancellationRequested) break; // Check positions that are near expiry, recently expired, or have no expiry but have a slug var checkPositions = acc.OpenPositions.Values.Where(p => !string.IsNullOrEmpty(p.MarketSlug) && ( // Has expiry and is within check window (-1 day to +30 days) (p.ExpiryDate.HasValue && (DateTime.UtcNow - p.ExpiryDate.Value).TotalDays > -1 && (DateTime.UtcNow - p.ExpiryDate.Value).TotalDays < 30) || // No expiry date at all — always check via API !p.ExpiryDate.HasValue )).ToList(); foreach (var pos in checkPositions) { var (isClosed, isWinner) = await _api.CheckMarketResolutionAsync(pos.MarketSlug, pos.TokenId); if (isClosed) { decimal exitPrice = isWinner ? 1.0m : 0.0m; _logger.Info($"🏆 Demo Market {pos.MarketQuestion} aufgelöst! Auszahlung: ${(exitPrice * pos.Size):F2}"); var signal = new CopySignal { TraderId = 0, TokenId = pos.TokenId, MarketSlug = pos.MarketSlug, MarketQuestion = pos.MarketQuestion, Outcome = pos.Outcome, Side = "SELL", Price = exitPrice, Size = pos.Size, Timestamp = DateTime.UtcNow, Reason = "Market Resolved" }; _signalWriter.TryWrite(signal); await Task.Delay(500, ct); } } } } private async Task PollLiveAccountsAsync(CancellationToken ct) { // Always sync live positions so the Dashboard UI accurately reflects open PnL and portfolio balance var liveAccounts = _state.Accounts.Values.Where(a => !a.IsDemo && a.IsActive && !string.IsNullOrEmpty(a.WalletAddress)).ToList(); if (liveAccounts.Count == 0) return; foreach (var acc in liveAccounts) { if (ct.IsCancellationRequested) break; var posList = await _api.SyncOpenPositionsAsync(acc.WalletAddress); if (posList.Count == 0) continue; var currentTokens = new HashSet(); foreach (var posJson in posList) { string asset = posJson.TryGetProperty("asset", out var ap) ? ap.GetString() ?? "" : ""; if (string.IsNullOrEmpty(asset)) continue; currentTokens.Add(asset); string slug = posJson.TryGetProperty("slug", out var sp) ? sp.GetString() ?? "" : ""; string title = posJson.TryGetProperty("title", out var tp) ? tp.GetString() ?? "" : ""; string opp = posJson.TryGetProperty("oppositeOutcome", out var op) ? op.GetString() ?? "" : "No"; decimal size = 0m, entryPrice = 0m, amountUsd = 0m, curPrice = 0m, curValue = 0m; if (posJson.TryGetProperty("size", out var sprop)) size = ParseDecimal(sprop); if (posJson.TryGetProperty("avgPrice", out var aprop)) entryPrice = ParseDecimal(aprop); // Critical Fix: "totalBought" is size. "initialValue" is original USD investment cost. if (posJson.TryGetProperty("initialValue", out var tbprop)) amountUsd = ParseDecimal(tbprop); if (posJson.TryGetProperty("curPrice", out var cpprop)) curPrice = ParseDecimal(cpprop); if (posJson.TryGetProperty("currentValue", out var cvprop)) curValue = ParseDecimal(cvprop); DateTime? expiry = null; if (posJson.TryGetProperty("endDate", out var ep)) { if (DateTime.TryParse(ep.GetString(), out var ed)) expiry = DateTime.SpecifyKind(ed.Date, DateTimeKind.Utc); } if (acc.OpenPositions.TryGetValue(asset, out var existing)) { existing.Size = size; existing.EntryPrice = entryPrice; existing.AmountUsd = amountUsd; existing.CurrentPrice = curPrice; existing.CurrentValueUsd = curValue; if (expiry.HasValue) existing.ExpiryDate = expiry; } else { var newPos = new Position { TokenId = asset, MarketSlug = slug, MarketQuestion = title, Outcome = opp == "Yes" ? "No" : "Yes", SourceTraderName = "Live Sync", Side = "BUY", Size = size, EntryPrice = entryPrice, AmountUsd = amountUsd, CurrentPrice = curPrice, CurrentValueUsd = curValue, ExpiryDate = expiry }; acc.OpenPositions.TryAdd(asset, newPos); _logger.Info($"🌐 Live Position erkannt: {title} ({newPos.Outcome}) - ${amountUsd} - Account: {acc.Name}"); } } var tokensToRemove = acc.OpenPositions .Where(kvp => !currentTokens.Contains(kvp.Key)) .Where(kvp => (DateTime.UtcNow - kvp.Value.OpenedAt).TotalMinutes > 5) .Select(kvp => kvp.Key) .ToList(); if (tokensToRemove.Count > 0) { var closedPositions = await _api.SyncClosedPositionsAsync(acc.WalletAddress, 50); foreach (var k in tokensToRemove) { if (acc.OpenPositions.TryRemove(k, out var removedPos)) { JsonElement? matchedClose = null; foreach (var cm in closedPositions) { if (cm.TryGetProperty("asset", out var ap) && ap.GetString() == k) { matchedClose = cm; break; } } if (matchedClose.HasValue) { decimal realizedPnl = 0m; if (matchedClose.Value.TryGetProperty("realizedPnl", out var rPnlProp)) realizedPnl = ParseDecimal(rPnlProp); _state.GlobalPnl += realizedPnl; decimal exitPrice = removedPos.Size > 0 ? (removedPos.AmountUsd + realizedPnl) / removedPos.Size : 0m; string duplicateKey = $"{acc.AccountId}_{removedPos.TokenId}"; if (!_processedClosures.ContainsKey(duplicateKey)) { _logger.Info($"🏆 Live Market {removedPos.MarketQuestion} geschlossen! PnL: ${(realizedPnl):F2}"); var ctRecord = new ClosedTrade { TradeId = _state.TotalCopyTrades, AccountId = acc.AccountId, SourceTraderId = removedPos.SourceTraderId, IsDemo = false, MarketSlug = removedPos.MarketSlug, MarketQuestion = removedPos.MarketQuestion, Outcome = removedPos.Outcome, Side = "SELL", EntryPrice = removedPos.EntryPrice, ExitPrice = exitPrice, Size = removedPos.Size, RealizedPnl = realizedPnl, PnlPercent = removedPos.AmountUsd > 0 ? (realizedPnl / removedPos.AmountUsd * 100m) : 0m, OpenedAt = removedPos.OpenedAt, ClosedAt = DateTime.UtcNow, ExitReason = "API Closed" }; _processedClosures.TryAdd(duplicateKey, true); _closedTradeWriter.TryWrite(ctRecord); } } else { var (isClosed, isWinner) = await _api.CheckMarketResolutionAsync(removedPos.MarketSlug, removedPos.TokenId); if (isClosed) { decimal exitPrice = isWinner ? 1.0m : 0.0m; decimal exitUsd = removedPos.Size * exitPrice; decimal realizedPnl = exitUsd - removedPos.AmountUsd; _state.GlobalPnl += realizedPnl; string duplicateKey = $"{acc.AccountId}_{removedPos.TokenId}"; if (!_processedClosures.ContainsKey(duplicateKey)) { _logger.Info($"🏆 Live Market {removedPos.MarketQuestion} aufgelöst (Fallback)! Auszahlung: ${(exitPrice * removedPos.Size):F2}"); var ctRecord = new ClosedTrade { TradeId = _state.TotalCopyTrades, AccountId = acc.AccountId, SourceTraderId = removedPos.SourceTraderId, IsDemo = false, MarketSlug = removedPos.MarketSlug, MarketQuestion = removedPos.MarketQuestion, Outcome = removedPos.Outcome, Side = "SELL", EntryPrice = removedPos.EntryPrice, ExitPrice = exitPrice, Size = removedPos.Size, RealizedPnl = realizedPnl, PnlPercent = removedPos.AmountUsd > 0 ? (realizedPnl / removedPos.AmountUsd * 100m) : 0m, OpenedAt = removedPos.OpenedAt, ClosedAt = DateTime.UtcNow, ExitReason = "API Resolved" }; _processedClosures.TryAdd(duplicateKey, true); _closedTradeWriter.TryWrite(ctRecord); } if (isWinner) { /* * DEATIVIERT: Automatischer Redeem via Python Script ist vorerst pausiert. * User kann die gewonnenen Shares per Klick im Polymarket Web-Interface redeemen. * Die Datenbank hat die PnL trotzdem bereits korrekt aufgezeichnet! * try { System.Diagnostics.Process.Start(new System.Diagnostics.ProcessStartInfo { FileName = "python", Arguments = $"redeem_markets.py {removedPos.TokenId} {acc.ApiKey} {acc.PrivateKey} {acc.ApiPassphrase}", UseShellExecute = false, CreateNoWindow = true }); _logger.Info($"Python Redeem Script für Token {removedPos.TokenId} asynchron ausgeführt."); } catch (Exception ex) { _logger.Error($"Fehler beim Starten von redeem_markets.py: {ex.Message}"); } */ _logger.Info($"🏆 Token {removedPos.TokenId} bereit für manuellen Redeem via Polymarket-Webseite. (P&L wurde bereits gebucht)."); } } else { _logger.Info($"🌐 Live Position {removedPos.MarketQuestion} (Ext. Verkauft/Wartend)"); } } } } } await Task.Delay(500, ct); } } private decimal ParseDecimal(JsonElement prop) { if (prop.ValueKind == JsonValueKind.Number) return prop.GetDecimal(); if (prop.ValueKind == JsonValueKind.String && decimal.TryParse(prop.GetString(), System.Globalization.NumberStyles.Any, System.Globalization.CultureInfo.InvariantCulture, out var parsed)) return parsed; return 0m; } private void ProcessActivityItem(JsonElement act, TrackedTrader trader) { try { string txHash = act.GetProperty("transactionHash").GetString() ?? ""; if (string.IsNullOrEmpty(txHash) || _processedTxHashes.ContainsKey(txHash)) return; // Duplicate or invalid string type = act.GetProperty("type").GetString() ?? ""; if (type.ToUpper() != "TRADE" && type.ToUpper() != "BUY" && type.ToUpper() != "SELL") return; string sideStr = type; // Fallback to type if (act.TryGetProperty("side", out var sideProp) && sideProp.ValueKind == JsonValueKind.String) sideStr = sideProp.GetString() ?? sideStr; else if (act.TryGetProperty("action", out var actionProp) && actionProp.ValueKind == JsonValueKind.String) sideStr = actionProp.GetString() ?? sideStr; else if (act.TryGetProperty("tradeType", out var ttProp) && ttProp.ValueKind == JsonValueKind.String) sideStr = ttProp.GetString() ?? sideStr; string asset = ""; if (act.TryGetProperty("asset", out var assetProp) && assetProp.ValueKind == JsonValueKind.String) asset = assetProp.GetString() ?? ""; if (string.IsNullOrEmpty(asset) && act.TryGetProperty("tokenId", out var tidProp) && tidProp.ValueKind == JsonValueKind.String) asset = tidProp.GetString() ?? ""; if (string.IsNullOrEmpty(asset) && act.TryGetProperty("token_id", out var t_idProp) && t_idProp.ValueKind == JsonValueKind.String) asset = t_idProp.GetString() ?? ""; if (string.IsNullOrEmpty(asset) && act.TryGetProperty("conditionId", out var cidProp) && cidProp.ValueKind == JsonValueKind.String) asset = cidProp.GetString() ?? ""; if (string.IsNullOrEmpty(asset) && act.TryGetProperty("condition_id", out var c_idProp) && c_idProp.ValueKind == JsonValueKind.String) asset = c_idProp.GetString() ?? ""; decimal price = 0m; if (act.TryGetProperty("price", out var priceProp)) { if (priceProp.ValueKind == JsonValueKind.Number) price = priceProp.GetDecimal(); else if (priceProp.ValueKind == JsonValueKind.String) decimal.TryParse(priceProp.GetString(), out price); } decimal size = 0m; if (act.TryGetProperty("size", out var sizeProp)) { if (sizeProp.ValueKind == JsonValueKind.Number) size = sizeProp.GetDecimal(); else if (sizeProp.ValueKind == JsonValueKind.String) decimal.TryParse(sizeProp.GetString(), out size); } // Parse timestamp to prevent old trades DateTime tradeTs = DateTime.UtcNow; if (act.TryGetProperty("timestamp", out var tsProp)) { if (tsProp.ValueKind == JsonValueKind.Number) // Unix tradeTs = DateTimeOffset.FromUnixTimeSeconds(tsProp.GetInt64()).UtcDateTime; else if (tsProp.ValueKind == JsonValueKind.String && DateTime.TryParse(tsProp.GetString(), out var dt)) tradeTs = dt.ToUniversalTime(); } // If trade is older than 120 seconds, skip if ((DateTime.UtcNow - tradeTs).TotalSeconds > 120) { // Still add to seen so we don't re-parse it _processedTxHashes.TryAdd(txHash, DateTime.UtcNow); return; } _processedTxHashes.TryAdd(txHash, DateTime.UtcNow); var displayQuestion = ""; if (act.TryGetProperty("title", out var titleProp)) displayQuestion = titleProp.GetString() ?? ""; var signal = new CopySignal { TraderId = trader.Id, TokenId = asset, ConditionId = "", MarketSlug = act.TryGetProperty("slug", out var sp) ? sp.GetString() ?? "" : (act.TryGetProperty("marketSlug", out var msp) ? msp.GetString() ?? "" : ""), Side = sideStr.ToUpper().Contains("SELL") ? "SELL" : "BUY", Price = price, Size = size, Timestamp = tradeTs, MarketQuestion = displayQuestion, Outcome = act.TryGetProperty("outcome", out var outProp) ? outProp.GetString() ?? "" : "", Reason = sideStr.ToUpper().Contains("SELL") ? "Master Trader Sold" : "" }; // Parse endDate from activity JSON for market expiry if (act.TryGetProperty("endDate", out var endDateProp)) { if (endDateProp.ValueKind == JsonValueKind.String && DateTime.TryParse(endDateProp.GetString(), null, System.Globalization.DateTimeStyles.RoundtripKind, out var endDt)) signal.EndDate = endDt.ToUniversalTime(); else if (endDateProp.ValueKind == JsonValueKind.Number) signal.EndDate = DateTimeOffset.FromUnixTimeSeconds(endDateProp.GetInt64()).UtcDateTime; } else if (act.TryGetProperty("end_date_iso", out var endIso) && endIso.ValueKind == JsonValueKind.String) { if (DateTime.TryParse(endIso.GetString(), null, System.Globalization.DateTimeStyles.RoundtripKind, out var endDt2)) signal.EndDate = endDt2.ToUniversalTime(); } string shareType = string.IsNullOrEmpty(signal.Outcome) ? signal.Side : signal.Outcome; _logger.Trade($"🚨 [QUELLE: {trader.DisplayName}] Neuer Trade erkannt!\n" + $" Markt: {signal.MarketQuestion}\n" + $" Aktion: {signal.Side} {shareType} ({signal.Size:F2} Shares @ ${signal.Price:F3})\n" + $" Zeit: {signal.Timestamp:HH:mm:ss} UTC"); // Push to the processing queue _signalWriter.TryWrite(signal); } catch (Exception ex) { _logger.Warning($"Fehler beim Parsen einer Activity JSON: {ex.Message}"); } } } }