Implement B4: Add Predictalytics.Application.Tests project with unit tests for PositionPnLEngine and AnalyticsService

This commit is contained in:
Richard
2026-07-03 11:22:43 +02:00
parent f8c8230d99
commit 7a44914d9d
7 changed files with 459 additions and 1 deletions
@@ -0,0 +1,174 @@
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;
/// <summary>
/// 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.
/// </summary>
public class TradeRetentionWorker : BackgroundService
{
private readonly IServiceProvider _services;
private readonly IConfiguration _config;
private readonly ILogger<TradeRetentionWorker> _logger;
private readonly TimeSpan _runInterval = TimeSpan.FromHours(24);
public TradeRetentionWorker(IServiceProvider services, IConfiguration config, ILogger<TradeRetentionWorker> 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 (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<AppDbContext>();
// Load configuration values
var retentionDays = _config.GetValue("RetentionSettings:RetentionDays", 90);
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);
// 1. Prune Old Trades (C1 & C2)
_logger.LogInformation("Pruning trades older than {Cutoff}...", retentionCutoff);
var deletedCount = await db.Trades
.Where(t => t.ExecutedAt < retentionCutoff)
.ExecuteDeleteAsync(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;
// 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_"))
.ToListAsync(ct);
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;
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
};
// Remove the individual trades
db.Trades.RemoveRange(list);
// Add the compacted trade
db.Trades.Add(compactedTrade);
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.");
}
}