diff --git a/src/Predictalytics.Api/Endpoints/JobEndpoints.cs b/src/Predictalytics.Api/Endpoints/JobEndpoints.cs index e3fe720..d5479aa 100644 --- a/src/Predictalytics.Api/Endpoints/JobEndpoints.cs +++ b/src/Predictalytics.Api/Endpoints/JobEndpoints.cs @@ -53,5 +53,33 @@ public static class JobEndpoints await repo.AddAsync(job, ct); return Results.Ok(job.Id); }); + + group.MapPost("/analyze-backlog", async (int? take, IJobRepository repo, Predictalytics.Infrastructure.Data.AppDbContext db, CancellationToken ct) => + { + int batchSize = take ?? 50; + + var traderIds = await Microsoft.EntityFrameworkCore.EntityFrameworkQueryableExtensions.ToListAsync( + db.Traders + .Where(t => t.LastAnalyzedAt == null || (t.LastTradesUpdatedAt != null && t.LastTradesUpdatedAt > t.LastAnalyzedAt)) + .OrderBy(t => t.LastAnalyzedAt == null ? 0 : 1) + .ThenByDescending(t => t.TotalTrades) + .Select(t => t.Id) + .Take(batchSize), + ct); + + int queued = 0; + foreach (var tId in traderIds) + { + var job = new BackgroundJob + { + JobType = JobType.TraderAnalysis, + Status = JobStatus.Pending, + TraderId = tId + }; + await repo.AddAsync(job, ct); + queued++; + } + return Results.Ok(new { Queued = queued }); + }); } } diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs index c2f63d2..1b54b1f 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs @@ -12,7 +12,13 @@ public class TraderRepository : ITraderRepository public TraderRepository(AppDbContext db) => _db = db; public async Task GetByIdAsync(int id, CancellationToken ct = default) - => await _db.Traders.Include(t => t.CurrentScore).Include(t => t.CategoryPerformances).FirstOrDefaultAsync(t => t.Id == id, ct); + { + return await _db.Traders + .Include(t => t.CurrentScore) + .Include(t => t.Analytics) + .Include(t => t.CategoryPerformances) + .FirstOrDefaultAsync(t => t.Id == id, ct); + } public async Task GetByPlatformIdAsync(PlatformType platform, string platformUserId, CancellationToken ct = default) => await _db.Traders.Include(t => t.CurrentScore).Include(t => t.CategoryPerformances) diff --git a/src/Predictalytics.WinFormsHost/MainForm.cs b/src/Predictalytics.WinFormsHost/MainForm.cs index fd7e6f0..9137c60 100644 --- a/src/Predictalytics.WinFormsHost/MainForm.cs +++ b/src/Predictalytics.WinFormsHost/MainForm.cs @@ -58,8 +58,6 @@ public partial class MainForm : Form // Wire up button events btn_serverstart.Click += Btn_serverstart_Click; btn_localWebserver.Click += Btn_localWebserver_Click; - btn_syncmarkets.Click += syncMarketsaToolStripMenuItem_Click; - btn_dbUpdate.Click += btn_dbUpdate_Click; Log.Information("MainForm initialized. Ready."); Log.Information("Press 'Start Server' to begin polling & discovery."); diff --git a/src/Predictalytics.WinFormsHost/Services/EmbeddedWebServer.cs b/src/Predictalytics.WinFormsHost/Services/EmbeddedWebServer.cs index 14a0c4a..3131380 100644 --- a/src/Predictalytics.WinFormsHost/Services/EmbeddedWebServer.cs +++ b/src/Predictalytics.WinFormsHost/Services/EmbeddedWebServer.cs @@ -181,16 +181,11 @@ public class EmbeddedWebServer await marketRepo.AddOrUpdateEventsAsync(events, ct); totalSynced += events.Count; - offset += batchSize; + offset += events.Count; if (totalSynced % 500 == 0 || events.Count < batchSize) Log.Information("[{Platform}] Synced {Total} events so far (offset={Offset}, includeClosed={Closed})...", provider.PlatformName, totalSynced, offset, includeClosed); - if (events.Count < batchSize) - { - Log.Warning("[{Platform}] Batch was smaller than limit ({Count}/{Limit}), assuming end of list.", provider.PlatformName, events.Count, batchSize); - break; - } } catch (Exception ex) { diff --git a/src/Predictalytics.Worker/Services/MarketHistoryWorker.cs b/src/Predictalytics.Worker/Services/MarketHistoryWorker.cs index 06000cb..a0991e1 100644 --- a/src/Predictalytics.Worker/Services/MarketHistoryWorker.cs +++ b/src/Predictalytics.Worker/Services/MarketHistoryWorker.cs @@ -5,6 +5,7 @@ using Predictalytics.Infrastructure.Logging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Configuration; namespace Predictalytics.Worker.Services; @@ -75,12 +76,15 @@ public class MarketHistoryWorker : BackgroundService .Distinct() .ToList(); + var config = scope.ServiceProvider.GetRequiredService(); + bool isEnabled = config.GetValue($"PlatformSettings:{market.Platform}:EnableCrawling", market.Platform == PlatformType.Polymarket); + foreach (var wallet in wallets) { var existing = await traderRepo.GetByPlatformIdAsync( market.Platform, wallet, stoppingToken); - if (existing == null) + if (existing == null && isEnabled) { var trader = new Domain.Entities.Trader { diff --git a/src/Predictalytics.Worker/Services/MarketSyncWorker.cs b/src/Predictalytics.Worker/Services/MarketSyncWorker.cs index 9d6e007..62f5023 100644 --- a/src/Predictalytics.Worker/Services/MarketSyncWorker.cs +++ b/src/Predictalytics.Worker/Services/MarketSyncWorker.cs @@ -80,7 +80,7 @@ public class MarketSyncWorker : BackgroundService passSynced += events.Count; cycleTotalSynced += events.Count; _statsService.TrackMarketSync(provider.Platform, events.Count); - offset += batchSize; + offset += events.Count; if (passSynced % 500 == 0) _logger.LogWarning("[{Platform}] Synced {Total} markets so far (includeClosed={Closed})...", p.PlatformName, passSynced, includeClosed); diff --git a/src/Predictalytics.Worker/Services/TopHolderDiscoveryWorker.cs b/src/Predictalytics.Worker/Services/TopHolderDiscoveryWorker.cs index 1301270..bc286fc 100644 --- a/src/Predictalytics.Worker/Services/TopHolderDiscoveryWorker.cs +++ b/src/Predictalytics.Worker/Services/TopHolderDiscoveryWorker.cs @@ -5,6 +5,7 @@ using Predictalytics.Infrastructure.Logging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Configuration; namespace Predictalytics.Worker.Services; @@ -36,15 +37,18 @@ public class TopHolderDiscoveryWorker : BackgroundService var providers = scope.ServiceProvider.GetRequiredService>(); var rateLimiter = scope.ServiceProvider.GetRequiredService(); - // Get top active markets by volume var activeMarkets = await marketRepo.GetActiveAsync(100, stoppingToken); _logger.LogInformation("👥 Scanning top holders across {Count} active markets", activeMarkets.Count); + var config = scope.ServiceProvider.GetRequiredService(); int totalDiscovered = 0; foreach (var market in activeMarkets) { if (stoppingToken.IsCancellationRequested) break; + + bool isEnabled = config.GetValue($"PlatformSettings:{market.Platform}:EnableCrawling", market.Platform == PlatformType.Polymarket); + if (!isEnabled) continue; var provider = providers.FirstOrDefault(p => p.Platform == market.Platform && p.IsImplemented); if (provider == null) continue; diff --git a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs index b96adfe..9a0c976 100644 --- a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs +++ b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs @@ -62,12 +62,17 @@ public class TraderAnalyticsWorker : BackgroundService else { var db = scope.ServiceProvider.GetRequiredService(); - // Find traders who have never been analyzed, or whose last analysis was before their latest trade. // Bug 5 Fix: Add 30-minute cooldown to prevent CPU looping. var cooldown = DateTime.UtcNow.AddMinutes(-30); + + // Find traders who have never been analyzed, or whose last analysis was before their latest trade. + // Include Resolution-Trigger: Traders with open positions in resolved markets. traderIds = await db.Traders - .Where(t => t.LastAnalyzedAt == null || - (t.LastAnalyzedAt < cooldown && t.Trades.Any(tr => tr.ExecutedAt > t.LastAnalyzedAt))) + .Where(t => + t.LastAnalyzedAt == null || + (t.LastAnalyzedAt < cooldown && t.LastTradesUpdatedAt != null && t.LastTradesUpdatedAt > t.LastAnalyzedAt) || + (t.LastAnalyzedAt < cooldown && t.Positions.Any(p => p.SharesHeld > 0 && p.MarketOutcome != null && p.MarketOutcome.Market != null && p.MarketOutcome.Market.IsResolved)) + ) .OrderBy(t => t.LastAnalyzedAt == null ? 0 : 1) .ThenBy(t => t.LastAnalyzedAt) .Select(t => t.Id) @@ -111,14 +116,21 @@ public class TraderAnalyticsWorker : BackgroundService var analyticsObj = trader.Analytics ?? new Predictalytics.Domain.Entities.TraderAnalytics { TraderId = trader.Id }; // Persist advanced copyability and quality scores derived from tape replay - analyticsObj.CopytradingScore = estScores.CopyabilityScore; + analyticsObj.CopytradingScore = estScores.CombinedScore; analyticsObj.CopytradingQualityScore = estScores.QualityScore; analyticsObj.CopytradingCopyabilityScore = estScores.CopyabilityScore; trader.Analytics = analyticsObj; + + // Only stamp if there were actually trades to analyze + trader.LastAnalyzedAt = DateTime.UtcNow; + await traderRepo.UpdateAsync(trader, ct); + } + else if (trader.LastAnalyzedAt == null) + { + // If it's a completely new trader with no trades, we still don't stamp LastAnalyzedAt + // so it remains null until the TradeHistoryWorker pulls trades. } - trader.LastAnalyzedAt = DateTime.UtcNow; - await traderRepo.UpdateAsync(trader, ct); var statsService = traderScope.ServiceProvider.GetService(); if (statsService != null)