using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Predictalytics.Domain.Entities; using Predictalytics.Domain.Enums; using Predictalytics.Infrastructure.Data; using System; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace Predictalytics.Worker.Services; /// /// Background service for storage optimization. /// Periodically deletes trades older than the retention window (default 90 days) /// and compacts high-frequency bot trades older than the compaction threshold (default 14 days) /// into daily summaries to reduce database row count. /// public class TradeRetentionWorker : BackgroundService { private readonly IServiceProvider _services; private readonly IConfiguration _config; private readonly ILogger _logger; private readonly TimeSpan _runInterval = TimeSpan.FromHours(24); public TradeRetentionWorker(IServiceProvider services, IConfiguration config, ILogger logger) { _services = services; _config = config; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("🧹 TradeRetentionWorker started"); await Task.Delay(15000, stoppingToken); // Let system initialize while (!stoppingToken.IsCancellationRequested) { try { await RunOptimizationAsync(stoppingToken); } catch (OperationCanceledException) { break; } catch (Exception ex) { _logger.LogError(ex, "Error occurred executing TradeRetentionWorker cycle."); } await Task.Delay(_runInterval, stoppingToken); } _logger.LogInformation("🧹 TradeRetentionWorker stopped"); } private async Task RunOptimizationAsync(CancellationToken ct) { using var scope = _services.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); // Load configuration values bool isEnabled = _config.GetValue("RetentionSettings:Enabled", true); if (!isEnabled) { _logger.LogInformation("🧹 TradeRetentionWorker: Retention is disabled in configuration. Skipping optimization."); return; } var retentionDays = _config.GetValue("RetentionSettings:RetentionDays", 180); var compactionDays = _config.GetValue("RetentionSettings:CompactionDays", 14); _logger.LogInformation("🧹 TradeRetentionWorker: Starting optimization. RetentionDays={Retention}, CompactionDays={Compaction}", retentionDays, compactionDays); var utcNow = DateTime.UtcNow; var retentionCutoff = utcNow.Date.AddDays(-retentionDays); var compactionCutoff = utcNow.Date.AddDays(-compactionDays); // Exclude trades if the trader is on any active Watchlist _logger.LogInformation("Pruning trades older than {Cutoff}...", retentionCutoff); var deletedTrades = await db.Trades .Where(t => t.ExecutedAt < retentionCutoff && !t.Trader.WatchlistEntries.Any() && t.Trader.IngestMode == IngestMode.Full) .Select(t => new { t.TraderId, t.MarketOutcomeId }) .Distinct() .ToListAsync(ct); var deletedCount = await db.Trades .Where(t => t.ExecutedAt < retentionCutoff && !t.Trader.WatchlistEntries.Any() && t.Trader.IngestMode == IngestMode.Full) .ExecuteDeleteAsync(ct); if (deletedCount > 0 && deletedTrades.Any()) { var validDeletedTrades = deletedTrades.Where(d => d.MarketOutcomeId.HasValue).ToList(); if (validDeletedTrades.Any()) { var posIdsToUpdate = await db.TraderPositions .Where(tp => validDeletedTrades.Select(d => d.TraderId).Contains(tp.TraderId) && validDeletedTrades.Select(d => d.MarketOutcomeId!.Value).Contains(tp.MarketOutcomeId)) .Select(tp => tp.Id) .ToListAsync(ct); // Filter on client side due to EF Core limitation with tuple Contains var actualPosIdsToUpdate = (await db.TraderPositions .Where(tp => posIdsToUpdate.Contains(tp.Id)) .ToListAsync(ct)) .Where(tp => validDeletedTrades.Any(d => d.TraderId == tp.TraderId && d.MarketOutcomeId == tp.MarketOutcomeId)) .Select(tp => tp.Id) .ToList(); if (actualPosIdsToUpdate.Any()) { await db.TraderPositions .Where(tp => actualPosIdsToUpdate.Contains(tp.Id)) .ExecuteUpdateAsync(s => s.SetProperty(p => p.IsHistoryPruned, true), ct); } } } _logger.LogInformation("Pruned {Count} old trades from the database.", deletedCount); // 2. Compact Bot Trades (C3) _logger.LogInformation("Beginning trade compaction for suspected bots/HF traders older than {Cutoff}...", compactionCutoff); var botTraderIds = await db.Traders .Where(t => t.IsSuspectedBot || t.Strategy == StrategyType.Bot) .Select(t => t.Id) .ToListAsync(ct); if (botTraderIds.Count == 0) { _logger.LogInformation("No bot/HF traders found for compaction."); return; } _logger.LogInformation("Found {Count} bot/HF traders to process.", botTraderIds.Count); foreach (var traderId in botTraderIds) { if (ct.IsCancellationRequested) break; // Fetch positions to ensure we only compact applied trades and can bump checkpoints var positions = await db.TraderPositions .Where(tp => tp.TraderId == traderId) .ToDictionaryAsync(tp => tp.MarketOutcomeId, ct); // Load candidate trades to compact (older than compactionCutoff, newer than retentionCutoff) var tradesToCompact = await db.Trades .Where(t => t.TraderId == traderId && t.ExecutedAt >= retentionCutoff && t.ExecutedAt < compactionCutoff && !t.PlatformTradeId.StartsWith("COMPACT_") && t.MarketOutcomeId != null) .ToListAsync(ct); // Filter strictly to trades that are already applied tradesToCompact = tradesToCompact .Where(t => positions.TryGetValue(t.MarketOutcomeId!.Value, out var pos) && t.Id <= pos.LastAppliedTradeId) .ToList(); if (tradesToCompact.Count == 0) continue; // Group trades by outcome, date, and side to aggregate var groups = tradesToCompact .GroupBy(t => new { t.MarketOutcomeId, Date = t.ExecutedAt.Date, t.Side }) .Where(g => g.Count() > 1 && g.Key.MarketOutcomeId.HasValue) .ToList(); if (groups.Count == 0) continue; _logger.LogInformation("Compacting {Count} groups of trades for trader {TraderId}...", groups.Count, traderId); int compactedTradeCount = 0; foreach (var g in groups) { var outcomeId = g.Key.MarketOutcomeId!.Value; var date = g.Key.Date; var side = g.Key.Side; var list = g.ToList(); var totalSize = list.Sum(t => t.Size); var totalAmount = list.Sum(t => t.Amount); if (totalSize <= 0) continue; if (positions.TryGetValue(outcomeId, out var pos)) { // Bug 2: Check for unapplied trades before compacting var hasUnappliedTrades = await db.Trades.AnyAsync( t => t.TraderId == traderId && t.MarketOutcomeId == outcomeId && t.Id > pos.LastAppliedTradeId, ct); if (hasUnappliedTrades) { continue; } } var weightedAvgPrice = totalAmount / totalSize; // Grab a representative trade to copy fields var sample = list.First(); var compactedTrade = new Trade { Platform = sample.Platform, PlatformTradeId = $"COMPACT_{traderId}_{outcomeId}_{date:yyyyMMdd}_{side}", TraderId = traderId, MarketId = sample.MarketId, DbMarketId = sample.DbMarketId, MarketOutcomeId = outcomeId, AssetId = sample.AssetId, Outcome = sample.Outcome, Side = side, Price = weightedAvgPrice, Size = totalSize, Amount = totalAmount, ExecutedAt = date.AddHours(12), // Set to noon of that day TransactionHash = null, AggregatedCount = list.Sum(t => t.AggregatedCount ?? 1) }; // Remove the individual trades db.Trades.RemoveRange(list); // Add the compacted trade db.Trades.Add(compactedTrade); // Save immediately so compactedTrade gets an ID await db.SaveChangesAsync(ct); // Bump the position checkpoint so it doesn't get double counted if (positions.TryGetValue(outcomeId, out var updatePos)) { updatePos.LastAppliedTradeId = Math.Max(updatePos.LastAppliedTradeId, compactedTrade.Id); // Mark as pruned so we don't accidentally reset and replay (which would lose the exact intraday timestamps) updatePos.IsHistoryPruned = true; db.TraderPositions.Update(updatePos); } compactedTradeCount += list.Count - 1; } if (compactedTradeCount > 0) { await db.SaveChangesAsync(ct); _logger.LogInformation("Reduced row count by {Count} for trader {TraderId}.", compactedTradeCount, traderId); } } _logger.LogInformation("🧹 TradeRetentionWorker: Optimization cycle complete."); } }