using System; using MongoDB.Driver; using PolyTraderSharp.Extensions; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Hosting; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services { public class MasterTraderAnalyticsJob : BackgroundService { private readonly TradingState _state; private readonly TerminalLogger _logger; private readonly IMongoDatabase _db; private readonly JobStatusRow _jobStatus; private readonly PolymarketApiService _api; public MasterTraderAnalyticsJob(TradingState state, TerminalLogger logger, IMongoDatabase db, JobManager jobManager, PolymarketApiService api) { _state = state; _logger = logger; _db = db; _api = api; _jobStatus = new JobStatusRow { JobName = "MasterTrader History", Description = "Überwacht die Performance aller Master-Trader (P&L, Winrate 7D).", StatusText = "Pending Initial Delay..." }; _jobStatus.ManualTriggerAction = async () => { _jobStatus.StatusText = "Running (Manual)..."; await RunHistoryAnalyticsAsync(); _jobStatus.StatusText = "Idle"; _jobStatus.LastRun = DateTime.Now; }; jobManager.RegisterJob(_jobStatus); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await Task.Delay(TimeSpan.FromSeconds(20), stoppingToken); // Start after other jobs while (!stoppingToken.IsCancellationRequested) { if (_jobStatus.IsEnabled) { try { _jobStatus.StatusText = "Running (Scheduled)..."; await RunHistoryAnalyticsAsync(); _jobStatus.LastRun = DateTime.Now; } catch (Exception ex) { _logger.Error($"Error in MasterTraderAnalyticsJob: {ex.Message}"); _jobStatus.StatusText = "Error!"; } finally { if (_jobStatus.StatusText != "Error!") _jobStatus.StatusText = "Idle"; } } else { _jobStatus.StatusText = "Paused"; } // Run twice a day (every 12 hours) _jobStatus.NextRun = DateTime.Now.AddHours(12); await Task.Delay(TimeSpan.FromHours(12), stoppingToken); } } public async Task RunHistoryAnalyticsAsync() { try { _logger.Info("🔄 Starte Master-Trader Historien-Download und Performance-Analyse..."); var historyColl = _db.GetCollection("mt_history"); historyColl.EnsureIndex(x => x.TraderId); historyColl.EnsureIndex(x => x.ClosedAt); DateTime cutoff7Days = DateTime.UtcNow.AddDays(-7); var tradersToAnalyze = _state.Traders.Values.Where(t => t.IsActive && !string.IsNullOrEmpty(t.WalletAddress)).ToList(); foreach (var trader in tradersToAnalyze) { try { // 1. Fetch History from Data API (100 is usually enough for 7 days) var closedPositions = await _api.SyncClosedPositionsAsync(trader.WalletAddress, 200); if (closedPositions.Count == 0) { continue; // Might be deleted or no history } int inserted = 0; foreach (var cp in closedPositions) { // parse timestamp DateTime closedTs = DateTime.UnixEpoch; if (cp.TryGetProperty("timestamp", out var tsProp)) { if (tsProp.ValueKind == JsonValueKind.Number) { long tsRaw = tsProp.GetInt64(); // if it's 13 digits (ms) vs 10 digits (s) if (tsRaw > 1000000000000) closedTs = DateTimeOffset.FromUnixTimeMilliseconds(tsRaw).UtcDateTime; else closedTs = DateTimeOffset.FromUnixTimeSeconds(tsRaw).UtcDateTime; } else if (tsProp.ValueKind == JsonValueKind.String && long.TryParse(tsProp.GetString(), out long tsStrRaw)) { if (tsStrRaw > 1000000000000) closedTs = DateTimeOffset.FromUnixTimeMilliseconds(tsStrRaw).UtcDateTime; else closedTs = DateTimeOffset.FromUnixTimeSeconds(tsStrRaw).UtcDateTime; } } // If trade is older than 14 days, ignore parsing to save DB space if (closedTs < DateTime.UtcNow.AddDays(-14)) continue; string tokenId = cp.TryGetProperty("asset", out var aProp) ? aProp.GetString() ?? "" : ""; decimal pnl = 0m; if (cp.TryGetProperty("realizedPnl", out var pProp)) { if (pProp.ValueKind == JsonValueKind.Number) pnl = pProp.GetDecimal(); else if (pProp.ValueKind == JsonValueKind.String && decimal.TryParse(pProp.GetString(), System.Globalization.NumberStyles.Any, System.Globalization.CultureInfo.InvariantCulture, out decimal nPnl)) { pnl = nPnl; } } // We can approximate uniqueness with TokenId & exact Time (+- 2 seconds) DateTime windowStart = closedTs.AddSeconds(-2); DateTime windowEnd = closedTs.AddSeconds(2); bool exists = historyColl.LiteFindOne(x => x.TraderId == trader.Id && x.TokenId == tokenId && x.ClosedAt >= windowStart && x.ClosedAt <= windowEnd) != null; if (!exists) { var record = new MasterTraderHistoryRecord { TraderId = trader.Id, TokenId = tokenId, ClosedAt = closedTs, RealizedPnl = pnl }; historyColl.Insert(record); inserted++; } } // Sleep to respect 10/s limits or general rate limits await Task.Delay(200); // 2. Calculate Stats from DB var last7DaysTrades = historyColl.LiteFind(x => x.TraderId == trader.Id && x.ClosedAt >= cutoff7Days).ToList(); trader.TotalTrades = last7DaysTrades.Count; trader.TotalPnl = (double)last7DaysTrades.Sum(x => x.RealizedPnl); // Treat positive PnL as win trader.WinningTrades = last7DaysTrades.Count(x => x.RealizedPnl > 0); trader.Winrate30t = trader.TotalTrades > 0 ? Math.Round(((double)trader.WinningTrades / trader.TotalTrades) * 100, 2) : 0; // Save updated trader to DB so UI updates var tColl = _db.GetCollection("tracked_traders"); tColl.Update(trader); if (inserted > 0 && trader.TotalTrades > 0) { _logger.Info($"📊 [MasterTrader: {trader.DisplayName}] - {inserted} neue Trades geladen. 7D: {trader.TotalTrades} Trades | PnL: ${trader.TotalPnl:F2} | Winrate: {trader.Winrate30t}%"); } } catch (Exception exInner) { _logger.Error($"Error processing history for MasterTrader {trader.DisplayName}: {exInner}"); } } _logger.Info("✅ Master-Trader Historien-Analyse abgeschlossen."); } catch (Exception ex) { _logger.Error($"MasterTraderAnalyticsJob Exception: {ex}"); } } } }