From 7a44914d9dab997571271e18dfe5e60e935f4aff Mon Sep 17 00:00:00 2001 From: Richard Date: Fri, 3 Jul 2026 11:22:43 +0200 Subject: [PATCH] Implement B4: Add Predictalytics.Application.Tests project with unit tests for PositionPnLEngine and AnalyticsService --- Predictalytics.slnx | 1 + .../Predictalytics.Application.Tests.csproj | 28 +++ .../Services/AnalyticsServiceTests.cs | 104 +++++++++++ .../Services/PositionPnLEngineTests.cs | 150 +++++++++++++++ .../Services/AnalyticsService.cs | 2 +- .../DependencyInjection.cs | 1 + .../Services/TradeRetentionWorker.cs | 174 ++++++++++++++++++ 7 files changed, 459 insertions(+), 1 deletion(-) create mode 100644 src/Predictalytics.Application.Tests/Predictalytics.Application.Tests.csproj create mode 100644 src/Predictalytics.Application.Tests/Services/AnalyticsServiceTests.cs create mode 100644 src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs create mode 100644 src/Predictalytics.Worker/Services/TradeRetentionWorker.cs diff --git a/Predictalytics.slnx b/Predictalytics.slnx index 5a926c5..95e1c7c 100644 --- a/Predictalytics.slnx +++ b/Predictalytics.slnx @@ -1,6 +1,7 @@ + diff --git a/src/Predictalytics.Application.Tests/Predictalytics.Application.Tests.csproj b/src/Predictalytics.Application.Tests/Predictalytics.Application.Tests.csproj new file mode 100644 index 0000000..0e19570 --- /dev/null +++ b/src/Predictalytics.Application.Tests/Predictalytics.Application.Tests.csproj @@ -0,0 +1,28 @@ + + + + net10.0 + enable + enable + false + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/Predictalytics.Application.Tests/Services/AnalyticsServiceTests.cs b/src/Predictalytics.Application.Tests/Services/AnalyticsServiceTests.cs new file mode 100644 index 0000000..4c0b791 --- /dev/null +++ b/src/Predictalytics.Application.Tests/Services/AnalyticsServiceTests.cs @@ -0,0 +1,104 @@ +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging.Abstractions; +using Predictalytics.Application.Interfaces; +using Predictalytics.Application.Services; +using Predictalytics.Domain.Entities; +using Predictalytics.Domain.Enums; +using Predictalytics.Domain.Interfaces; +using Predictalytics.Infrastructure.Data; +using Predictalytics.Infrastructure.Data.Repositories; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Xunit; + +namespace Predictalytics.Application.Tests.Services; + +public class AnalyticsServiceTests +{ + private AppDbContext CreateDbContext() + { + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase(databaseName: Guid.NewGuid().ToString()) + .Options; + return new AppDbContext(options); + } + + [Fact] + public async Task GetTraderDeepDiveAsync_CalculatesCorrectMetrics() + { + // Arrange + using var db = CreateDbContext(); + var traderRepo = new TraderRepository(db); + var tradeRepo = new TradeRepository(db); + var marketRepo = new MarketRepository(db); + + var discoveryMock = new MockDiscoveryService(); + var providers = new List(); + + var analyticsService = new AnalyticsService( + traderRepo, + tradeRepo, + null!, // alertRepo (not used in GetTraderDeepDiveAsync) + null!, // watchlistRepo (not used in GetTraderDeepDiveAsync) + marketRepo, + discoveryMock, + providers, + NullLogger.Instance + ); + + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + db.Traders.Add(trader); + + var market = new Market { Id = 10, PlatformMarketId = "pm1", Question = "Q?" }; + var outcome = new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 0.50m }; + market.Outcomes.Add(outcome); + db.Markets.Add(market); + + var baseTime = DateTime.UtcNow.AddDays(-5); + db.Trades.Add(new Trade + { + Id = 501, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Buy, Price = 0.40m, Size = 100m, Amount = 40m, + ExecutedAt = baseTime + }); + + db.Trades.Add(new Trade + { + Id = 502, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Sell, Price = 0.60m, Size = 100m, Amount = 60m, + ExecutedAt = baseTime.AddHours(24) // 24 hours holding duration + }); + + db.MarketOutcomePriceSnapshots.Add(new MarketOutcomePriceSnapshot + { + Id = 1, MarketOutcomeId = 100, Price = 0.50m, Timestamp = baseTime.AddHours(2) + }); + db.MarketOutcomePriceSnapshots.Add(new MarketOutcomePriceSnapshot + { + Id = 2, MarketOutcomeId = 100, Price = 0.60m, Timestamp = baseTime.AddHours(12) + }); + + await db.SaveChangesAsync(); + + // Act + var deepDive = await analyticsService.GetTraderDeepDiveAsync(1); + + // Assert + Assert.NotNull(deepDive); + Assert.Equal(24.0, (double)deepDive.AvgHoldDurationHours, 2); + + // Entry Quality: Buy at 0.40, subsequent prices are 0.50 and 0.60 (avg 0.55). + // Entry Quality = 50 + ((0.55 - 0.40) / 0.40) * 100 = 50 + 0.375 * 100 = 87.5 + Assert.Equal(87.5m, deepDive.EntryQuality); + } + + private class MockDiscoveryService : IDiscoveryService + { + public Task ImportTraderAsync(PlatformType platform, string platformUserId, string displayName, CancellationToken ct) => Task.FromResult(0); + public Task ScanTopHoldersAsync(CancellationToken ct) => Task.CompletedTask; + public Task> RunDiscoveryAsync(PlatformType platform, CancellationToken ct) => Task.FromResult>(new List()); + } +} diff --git a/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs b/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs new file mode 100644 index 0000000..ab3db45 --- /dev/null +++ b/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs @@ -0,0 +1,150 @@ +using Microsoft.EntityFrameworkCore; +using Predictalytics.Domain.Entities; +using Predictalytics.Domain.Enums; +using Predictalytics.Infrastructure.Data; +using Predictalytics.Infrastructure.Services; +using Microsoft.Extensions.Logging.Abstractions; +using System; +using System.Threading.Tasks; +using Xunit; + +namespace Predictalytics.Application.Tests.Services; + +public class PositionPnLEngineTests +{ + private AppDbContext CreateDbContext() + { + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase(databaseName: Guid.NewGuid().ToString()) + .Options; + return new AppDbContext(options); + } + + [Fact] + public async Task RecalculateTraderPositionsAsync_BuyTrade_CreatesTraderPosition() + { + // Arrange + using var db = CreateDbContext(); + var pnlEngine = new PositionPnLEngine(db, NullLogger.Instance); + + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + var market = new Market { Id = 10, PlatformMarketId = "pm1", Question = "Q?" }; + var outcome = new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 0.60m }; + market.Outcomes.Add(outcome); + + db.Traders.Add(trader); + db.Markets.Add(market); + + var buyTrade = new Trade + { + Id = 500, + TraderId = 1, + DbMarketId = 10, + MarketOutcomeId = 100, + Side = TradeSide.Buy, + Price = 0.50m, + Size = 100m, + Amount = 50m, + ExecutedAt = DateTime.UtcNow + }; + db.Trades.Add(buyTrade); + await db.SaveChangesAsync(); + + // Act + await pnlEngine.RecalculateTraderPositionsAsync(1); + + // Assert + var pos = await db.TraderPositions.FirstOrDefaultAsync(p => p.TraderId == 1 && p.MarketOutcomeId == 100); + Assert.NotNull(pos); + Assert.Equal(100m, pos.SharesHeld); + Assert.Equal(0.50m, pos.AvgCost); + Assert.Equal(0m, pos.RealizedPnl); + } + + [Fact] + public async Task RecalculateTraderPositionsAsync_SellTrade_CalculatesRealizedPnL() + { + // Arrange + using var db = CreateDbContext(); + var pnlEngine = new PositionPnLEngine(db, NullLogger.Instance); + + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + var market = new Market { Id = 10, PlatformMarketId = "pm1", Question = "Q?" }; + var outcome = new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 0.60m }; + market.Outcomes.Add(outcome); + + db.Traders.Add(trader); + db.Markets.Add(market); + + // Buy 100 shares at 0.40 + db.Trades.Add(new Trade + { + Id = 501, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Buy, Price = 0.40m, Size = 100m, Amount = 40m, + ExecutedAt = DateTime.UtcNow.AddMinutes(-10) + }); + + // Sell 50 shares at 0.60 + db.Trades.Add(new Trade + { + Id = 502, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Sell, Price = 0.60m, Size = 50m, Amount = 30m, + ExecutedAt = DateTime.UtcNow + }); + + await db.SaveChangesAsync(); + + // Act + await pnlEngine.RecalculateTraderPositionsAsync(1); + + // Assert + var pos = await db.TraderPositions.FirstOrDefaultAsync(p => p.TraderId == 1 && p.MarketOutcomeId == 100); + Assert.NotNull(pos); + Assert.Equal(50m, pos.SharesHeld); + Assert.Equal(0.40m, pos.AvgCost); // AvgCost stays at 0.40 + Assert.Equal(10m, pos.RealizedPnl); // 50 * (0.60 - 0.40) = 10 + } + + [Fact] + public async Task RecalculateTraderPositionsAsync_RedeemTrade_CalculatesRedeemPnL() + { + // Arrange + using var db = CreateDbContext(); + var pnlEngine = new PositionPnLEngine(db, NullLogger.Instance); + + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + var market = new Market { Id = 10, PlatformMarketId = "pm1", Question = "Q?", IsResolved = true, ResolutionOutcome = "Yes" }; + var outcome = new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 1.00m }; + market.Outcomes.Add(outcome); + + db.Traders.Add(trader); + db.Markets.Add(market); + + // Buy 100 shares at 0.40 + db.Trades.Add(new Trade + { + Id = 501, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Buy, Price = 0.40m, Size = 100m, Amount = 40m, + ExecutedAt = DateTime.UtcNow.AddMinutes(-10) + }); + + // Redeem at 1.00 (Win) + db.Trades.Add(new Trade + { + Id = 503, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Redeem, Price = 1.00m, Size = 100m, Amount = 100m, + ExecutedAt = DateTime.UtcNow + }); + + await db.SaveChangesAsync(); + + // Act + await pnlEngine.RecalculateTraderPositionsAsync(1); + + // Assert + var pos = await db.TraderPositions.FirstOrDefaultAsync(p => p.TraderId == 1 && p.MarketOutcomeId == 100); + Assert.NotNull(pos); + Assert.Equal(0m, pos.SharesHeld); + Assert.Equal(60m, pos.RealizedPnl); // 100 * (1.00 - 0.40) = 60 + } +} diff --git a/src/Predictalytics.Application/Services/AnalyticsService.cs b/src/Predictalytics.Application/Services/AnalyticsService.cs index 1d646b1..007e633 100644 --- a/src/Predictalytics.Application/Services/AnalyticsService.cs +++ b/src/Predictalytics.Application/Services/AnalyticsService.cs @@ -285,7 +285,7 @@ public class AnalyticsService : IAnalyticsService botIndicators.Count > 0 ? StrategyType.Bot : StrategyType.Unknown; // Calculate holding duration - var holdDuration = (decimal)CalculateAvgHoldDuration(trades.ToList()); + var holdDuration = (decimal)CalculateAvgHoldDuration(trades.OrderBy(t => t.ExecutedAt).ToList()); // Calculate Entry and Exit qualities var entryQualities = new List(); diff --git a/src/Predictalytics.Worker/DependencyInjection.cs b/src/Predictalytics.Worker/DependencyInjection.cs index 1683ae7..e04a224 100644 --- a/src/Predictalytics.Worker/DependencyInjection.cs +++ b/src/Predictalytics.Worker/DependencyInjection.cs @@ -18,6 +18,7 @@ public static class DependencyInjection services.AddHostedService(); services.AddHostedService(); services.AddHostedService(); + services.AddHostedService(); return services; } } diff --git a/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs b/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs new file mode 100644 index 0000000..d6da9d4 --- /dev/null +++ b/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs @@ -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; + +/// +/// 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 (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 + 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."); + } +}