diff --git a/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs b/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs index 12f76bb..68d9360 100644 --- a/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs +++ b/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs @@ -36,6 +36,12 @@ public static class TraderEndpoints return Results.Ok(); }); + group.MapPost("/{id:int}/force-analyze", async (int id, IAnalyticsService svc, CancellationToken ct) => + { + await svc.ForceAnalyzeTraderAsync(id, ct); + return Results.Ok(); + }); + group.MapPost("/{id:int}/watchlist", async (int id, WatchlistService svc, CancellationToken ct) => { await svc.AddAsync(id, "Watched via UI", null, ct); diff --git a/src/Predictalytics.Application/Interfaces/IAnalyticsService.cs b/src/Predictalytics.Application/Interfaces/IAnalyticsService.cs index acba5e8..975456e 100644 --- a/src/Predictalytics.Application/Interfaces/IAnalyticsService.cs +++ b/src/Predictalytics.Application/Interfaces/IAnalyticsService.cs @@ -29,6 +29,9 @@ public interface IAnalyticsService /// Manually trigger a trade history sync for a specific trader. Task TriggerTradeSyncAsync(int traderId, CancellationToken ct = default); + /// Manually trigger a deep-dive analysis (PnL recalculation) for a specific trader. + Task ForceAnalyzeTraderAsync(int traderId, CancellationToken ct = default); + /// Manually add a trader by platform and wallet address. Task AddTraderAsync(string platform, string walletAddress, CancellationToken ct = default); } diff --git a/src/Predictalytics.Application/Services/AnalyticsService.cs b/src/Predictalytics.Application/Services/AnalyticsService.cs index e4e413e..1fc7f00 100644 --- a/src/Predictalytics.Application/Services/AnalyticsService.cs +++ b/src/Predictalytics.Application/Services/AnalyticsService.cs @@ -17,15 +17,22 @@ public class AnalyticsService : IAnalyticsService private readonly IMarketRepository _marketRepo; private readonly IDiscoveryService _discovery; private readonly IEnumerable _providers; + private readonly IPositionPnLEngine _pnlEngine; private readonly ILogger _logger; public AnalyticsService(ITraderRepository traderRepo, ITradeRepository tradeRepo, IAlertRepository alertRepo, IWatchlistRepository watchlistRepo, IMarketRepository marketRepo, - IDiscoveryService discovery, IEnumerable providers, ILogger logger) + IDiscoveryService discovery, IEnumerable providers, IPositionPnLEngine pnlEngine, ILogger logger) { - _traderRepo = traderRepo; _tradeRepo = tradeRepo; - _alertRepo = alertRepo; _watchlistRepo = watchlistRepo; _marketRepo = marketRepo; - _discovery = discovery; _providers = providers; _logger = logger; + _traderRepo = traderRepo; + _tradeRepo = tradeRepo; + _alertRepo = alertRepo; + _watchlistRepo = watchlistRepo; + _marketRepo = marketRepo; + _discovery = discovery; + _providers = providers; + _pnlEngine = pnlEngine; + _logger = logger; } public async Task GetDashboardAsync(CancellationToken ct = default) @@ -419,6 +426,12 @@ public class AnalyticsService : IAnalyticsService } } + public async Task ForceAnalyzeTraderAsync(int traderId, CancellationToken ct = default) + { + _logger.LogInformation("Manually forcing deep analysis for trader {TraderId}", traderId); + await _pnlEngine.RecalculateTraderPositionsAsync(traderId, ct); + } + public async Task AddTraderAsync(string platform, string walletAddress, CancellationToken ct = default) { if (!Enum.TryParse(platform, true, out var pType)) diff --git a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs index 5671a2f..bdbf303 100644 --- a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs +++ b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs @@ -224,11 +224,12 @@ public class PositionPnLEngine : IPositionPnLEngine // Sync back to Trader record for quick sorting / UI display trader.TotalPnl = overallPnl; trader.WinRate = winRateOverall; + trader.LastAnalyzedAt = DateTime.UtcNow; // Calculate Category Performance var existingCatPerf = await _db.TraderCategoryPerformances .Where(tcp => tcp.TraderId == traderId) - .ToDictionaryAsync(tcp => tcp.Category, ct); + .ToDictionaryAsync(tcp => (tcp.Category, tcp.Subcategory), ct); var newCatPerf = CalculateCategoryPerformances(trades, tempPositions); @@ -348,11 +349,11 @@ public class PositionPnLEngine : IPositionPnLEngine return false; } - private static Dictionary CalculateCategoryPerformances( + private static Dictionary<(MarketCategory, string), TraderCategoryPerformance> CalculateCategoryPerformances( List trades, Dictionary finalPositions) { - var result = new Dictionary(); + var result = new Dictionary<(MarketCategory, string), TraderCategoryPerformance>(); var tradesByMarket = trades .Where(t => t.MarketOutcome?.Market != null) @@ -362,11 +363,14 @@ public class PositionPnLEngine : IPositionPnLEngine { var market = marketGroup.Key; var category = market.Category; + var subcat = market.Subcategory ?? ""; - if (!result.TryGetValue(category, out var perf)) + var key = (category, subcat); + + if (!result.TryGetValue(key, out var perf)) { - perf = new TraderCategoryPerformance { Category = category }; - result[category] = perf; + perf = new TraderCategoryPerformance { Category = category, Subcategory = subcat }; + result[key] = perf; } // Add volume diff --git a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs index 5b95ff1..ffb8dbb 100644 --- a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs +++ b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs @@ -40,21 +40,23 @@ public class TraderAnalyticsWorker : BackgroundService private async Task RunAnalyticsAsync(CancellationToken ct) { - var cutoff30d = DateTime.UtcNow.AddDays(-30); List traderIds; using (var scope = _services.CreateScope()) { var db = scope.ServiceProvider.GetRequiredService(); - // Find traders active in the last 30 days - traderIds = await db.Trades - .Where(t => t.ExecutedAt >= cutoff30d) - .Select(t => t.TraderId) - .Distinct() + // Find traders who have never been analyzed, or whose last analysis was before their latest trade. + // Prioritize never-analyzed traders. + traderIds = await db.Traders + .Where(t => t.LastAnalyzedAt == null || t.Trades.Any(tr => tr.ExecutedAt > t.LastAnalyzedAt)) + .OrderBy(t => t.LastAnalyzedAt == null ? 0 : 1) + .ThenBy(t => t.LastAnalyzedAt) + .Select(t => t.Id) + .Take(500) // Limit batch size to prevent long-running loops without save .ToListAsync(ct); } - _logger.LogInformation("Found {Count} active traders to analyze", traderIds.Count); + _logger.LogInformation("Found {Count} active or unanalyzed traders to update", traderIds.Count); foreach (var id in traderIds) {