diff --git a/src/Predictalytics.Application.Tests/Services/IngestModeTests.cs b/src/Predictalytics.Application.Tests/Services/IngestModeTests.cs new file mode 100644 index 0000000..0f25e7d --- /dev/null +++ b/src/Predictalytics.Application.Tests/Services/IngestModeTests.cs @@ -0,0 +1,227 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Data.Sqlite; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Predictalytics.Application.Interfaces; +using Predictalytics.Domain.Entities; +using Predictalytics.Domain.Enums; +using Predictalytics.Domain.Interfaces; +using Predictalytics.Infrastructure.Data; +using Predictalytics.Infrastructure.Data.Repositories; +using Predictalytics.Infrastructure.Services; +using Predictalytics.Worker.Services; +using Moq; +using Xunit; + +namespace Predictalytics.Application.Tests.Services; + +public class IngestModeTests +{ + [Fact] + public async Task PollingWorker_BypassesSnapshotOnlyTrader() + { + using var connection = new SqliteConnection("DataSource=:memory:"); + connection.Open(); + var options = new DbContextOptionsBuilder() + .UseSqlite(connection) + .Options; + + using (var setup = new AppDbContext(options)) + { + setup.Database.EnsureCreated(); + setup.Traders.Add(new Trader + { + Id = 1, + PlatformUserId = "0x1", + DisplayName = "HN1", + Platform = PlatformType.Polymarket, + IngestMode = IngestMode.SnapshotOnly + }); + setup.SaveChanges(); + } + + var services = new ServiceCollection(); + services.AddScoped(_ => new AppDbContext(options)); + services.AddScoped(); + + var mockTradeRepo = new Mock(); + mockTradeRepo.Setup(r => r.GetKnownPlatformTradeIdsAsync(It.IsAny(), It.IsAny(), It.IsAny>(), It.IsAny())) + .ReturnsAsync(new HashSet()); + services.AddSingleton(mockTradeRepo.Object); + + services.AddScoped(); + services.AddSingleton(typeof(ILogger<>), typeof(NullLogger<>)); + + var mockProvider = new Mock(); + mockProvider.Setup(p => p.Platform).Returns(PlatformType.Polymarket); + mockProvider.Setup(p => p.IsImplemented).Returns(true); + mockProvider.Setup(p => p.PlatformName).Returns("Polymarket"); + + mockProvider.Setup(p => p.GetTraderTradesAsync(It.IsAny(), It.IsAny(), It.IsAny())) + .ReturnsAsync(new List { new Trade { PlatformTradeId = "test_tx" } }); + + services.AddSingleton>(new[] { mockProvider.Object }); + + var mockRateLimiter = new Mock(); + services.AddSingleton(mockRateLimiter.Object); + + var mockStats = new Mock(); + services.AddSingleton(mockStats.Object); + + using var provider = services.BuildServiceProvider(); + + var worker = new PollingWorker(provider, NullLogger.Instance, mockStats.Object); + + var method = typeof(PollingWorker).GetMethod("ExecuteAsync", System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic); + Assert.NotNull(method); + + using var cts = new CancellationTokenSource(); + var task = (Task)method!.Invoke(worker, new object[] { cts.Token })!; + + // Wait for worker to start, process, and then cancel + await Task.Delay(6000); + cts.Cancel(); + + try { await task; } catch (Exception) { } + + // Assert: mockProvider GetTraderTradesAsync should NEVER be called since SnapshotOnly is skipped + mockProvider.Verify(p => p.GetTraderTradesAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); + } + + [Fact] + public async Task AggregatedMode_AggregatesTradesCorrectly() + { + using var connection = new SqliteConnection("DataSource=:memory:"); + connection.Open(); + var options = new DbContextOptionsBuilder() + .UseSqlite(connection) + .Options; + + using (var setup = new AppDbContext(options)) + { + setup.Database.EnsureCreated(); + setup.Traders.Add(new Trader + { + Id = 2, + PlatformUserId = "0x2", + DisplayName = "AggregatedTrader", + Platform = PlatformType.Polymarket, + IngestMode = IngestMode.Aggregated + }); + + var ev = new Event { Id = 1, Platform = PlatformType.Polymarket, Slug = "e", Title = "E" }; + setup.Set().Add(ev); + + var market = new Market { Id = 10, EventId = 1, PlatformMarketId = 1L, Question = "Q?" }; + market.Outcomes.Add(new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 0.45m }); + setup.Markets.Add(market); + + setup.SaveChanges(); + } + + var services = new ServiceCollection(); + services.AddScoped(_ => new AppDbContext(options)); + services.AddScoped(); + + var mockTradeRepo = new Mock(); + mockTradeRepo.Setup(r => r.GetKnownPlatformTradeIdsAsync(It.IsAny(), It.IsAny(), It.IsAny>(), It.IsAny())) + .ReturnsAsync(new HashSet()); + + // Mock AddRangeAsync to bypass MySql raw query and save trades directly to SQLite in-memory DB + mockTradeRepo.Setup(r => r.AddRangeAsync(It.IsAny>(), It.IsAny())) + .Callback, CancellationToken>((trades, ct) => + { + using var db = new AppDbContext(options); + foreach (var t in trades) + { + db.Trades.Add(t); + } + db.SaveChanges(); + }) + .Returns(Task.CompletedTask); + + services.AddSingleton(mockTradeRepo.Object); + + services.AddScoped(); + services.AddSingleton(typeof(ILogger<>), typeof(NullLogger<>)); + + var baseTime = new DateTime(2026, 7, 19, 10, 30, 0, DateTimeKind.Utc); + + var mockProvider = new Mock(); + mockProvider.Setup(p => p.Platform).Returns(PlatformType.Polymarket); + mockProvider.Setup(p => p.IsImplemented).Returns(true); + mockProvider.Setup(p => p.PlatformName).Returns("Polymarket"); + + // Fetch returns two trades in the same hour + mockProvider.Setup(p => p.GetTraderTradesAsync(It.IsAny(), It.IsAny(), It.IsAny())) + .ReturnsAsync(new List + { + new Trade + { + PlatformTradeId = "tx1", + MarketId = "10", + AssetId = "t100", + Side = TradeSide.Buy, + Price = 0.40m, + Size = 100m, + Amount = 40m, + ExecutedAt = baseTime + }, + new Trade + { + PlatformTradeId = "tx2", + MarketId = "10", + AssetId = "t100", + Side = TradeSide.Buy, + Price = 0.60m, + Size = 100m, + Amount = 60m, + ExecutedAt = baseTime.AddMinutes(15) + } + }); + + services.AddSingleton>(new[] { mockProvider.Object }); + + var mockRateLimiter = new Mock(); + services.AddSingleton(mockRateLimiter.Object); + + var mockStats = new Mock(); + services.AddSingleton(mockStats.Object); + + using var provider = services.BuildServiceProvider(); + + var worker = new PollingWorker(provider, NullLogger.Instance, mockStats.Object); + + var method = typeof(PollingWorker).GetMethod("ExecuteAsync", System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic); + Assert.NotNull(method); + + using var cts = new CancellationTokenSource(); + var task = (Task)method!.Invoke(worker, new object[] { cts.Token })!; + + // Wait for worker to start, process, and then cancel + await Task.Delay(6000); + cts.Cancel(); + + try { await task; } catch (Exception) { } + + // Verify that only 1 aggregated trade is saved in the database + using (var assertCtx = new AppDbContext(options)) + { + var trades = await assertCtx.Trades.ToListAsync(); + Assert.Single(trades); + var aggTrade = trades[0]; + Assert.Equal("AGG_2_100_Buy_2026071910", aggTrade.PlatformTradeId); + Assert.Equal(200m, aggTrade.Size); + Assert.Equal(100m, aggTrade.Amount); + Assert.Equal(0.50m, aggTrade.Price); // VWAP: (40 + 60) / (100 + 100) = 0.50 + Assert.Equal(2, aggTrade.AggregatedCount); + } + } +} diff --git a/src/Predictalytics.Application/Services/TraderTraitCalculator.cs b/src/Predictalytics.Application/Services/TraderTraitCalculator.cs index 43299c1..7985689 100644 --- a/src/Predictalytics.Application/Services/TraderTraitCalculator.cs +++ b/src/Predictalytics.Application/Services/TraderTraitCalculator.cs @@ -11,6 +11,11 @@ public static class TraderTraitCalculator IReadOnlyCollection positions) { var traits = new List<(string, decimal)>(); + if (trader.IngestMode == IngestMode.SnapshotOnly) + { + traits.Add(("not_copyable_hf", 1.0m)); + } + if (trades.Count == 0) return traits; var now = DateTime.UtcNow; diff --git a/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs b/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs index 981a61a..0943e74 100644 --- a/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs +++ b/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs @@ -58,7 +58,8 @@ public record TraderPositionInfo( decimal AveragePrice, decimal CurrentValue, decimal PnlPercent, - string? AssetId = null + string? AssetId = null, + decimal RealizedPnl = 0 ); /// diff --git a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs index 1a05433..6371a08 100644 --- a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs +++ b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs @@ -151,7 +151,7 @@ public class PolymarketProvider : IPlatformProvider platformUserId, r.Market, r.Question, r.Outcome, (decimal)r.Size, (decimal)r.AvgPrice, (decimal)r.CurrentValue, (decimal)r.PercentPnl, - r.AssetId + r.AssetId, (decimal)r.CashPnl )).ToList(); } diff --git a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs index 9c64f84..f2dd6ad 100644 --- a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs +++ b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs @@ -35,6 +35,12 @@ public class PositionPnLEngine : IPositionPnLEngine return; } + if (trader.IngestMode == IngestMode.SnapshotOnly) + { + _logger.LogInformation("Skipping tape replay/recalculation for SnapshotOnly trader {TraderId}", traderId); + return; + } + // Fetch all trades for this trader, sorted chronologically, including outcomes and markets var trades = await _db.Trades .Include(t => t.MarketOutcome) diff --git a/src/Predictalytics.Worker/Services/PollingWorker.cs b/src/Predictalytics.Worker/Services/PollingWorker.cs index 8e1c223..7462c4a 100644 --- a/src/Predictalytics.Worker/Services/PollingWorker.cs +++ b/src/Predictalytics.Worker/Services/PollingWorker.cs @@ -53,7 +53,7 @@ public class PollingWorker : BackgroundService var rateLimiter = scope.ServiceProvider.GetRequiredService(); var trader = await traderRepo.GetByIdAsync(t.Id, stoppingToken); - if (trader == null) continue; + if (trader == null || trader.IngestMode == IngestMode.SnapshotOnly) continue; var provider = providers.FirstOrDefault(p => p.Platform == trader.Platform && p.IsImplemented); if (provider == null) continue; @@ -122,6 +122,37 @@ public class PollingWorker : BackgroundService } } + if (trader.IngestMode == IngestMode.Aggregated && newTrades.Count > 0) + { + var aggregated = new List(); + foreach (var grp in newTrades.GroupBy(t => new { t.MarketOutcomeId, t.Side, Hour = t.ExecutedAt.ToString("yyyyMMddHH") })) + { + var first = grp.First(); + var totalAmount = grp.Sum(t => t.Amount); + var totalSize = grp.Sum(t => t.Size); + var vwap = totalSize > 0 ? totalAmount / totalSize : first.Price; + + var aggTrade = new Domain.Entities.Trade + { + PlatformTradeId = $"AGG_{trader.Id}_{grp.Key.MarketOutcomeId}_{grp.Key.Side}_{grp.Key.Hour}", + TraderId = trader.Id, + MarketOutcomeId = first.MarketOutcomeId, + DbMarketId = first.DbMarketId, + MarketId = first.MarketId, + AssetId = first.AssetId, + Outcome = first.Outcome, + Side = first.Side, + Price = vwap, + Amount = totalAmount, + Size = totalSize, + ExecutedAt = first.ExecutedAt, + AggregatedCount = grp.Count() + }; + aggregated.Add(aggTrade); + } + newTrades = aggregated; + } + // ── Persist new trades ── if (newTrades.Count > 0) { diff --git a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs index adf3b41..7490146 100644 --- a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs @@ -6,6 +6,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Configuration; +using Microsoft.EntityFrameworkCore; namespace Predictalytics.Worker.Services; @@ -173,7 +174,7 @@ public class TradeHistoryWorker : BackgroundService if (tradesPerDay > 5000) newMode = IngestMode.SnapshotOnly; else if (tradesPerDay > 100 && trader.IngestMode == IngestMode.Full) newMode = IngestMode.Aggregated; else if (trader.IngestMode == IngestMode.SnapshotOnly && tradesPerDay < 2500) newMode = IngestMode.Aggregated; - else if (trader.IngestMode == IngestMode.Aggregated && tradesPerDay < 50) newMode = IngestMode.Full; + else if (trader.IngestMode == IngestMode.Aggregated && tradesPerDay < 50 && days >= 7) newMode = IngestMode.Full; if (newMode != trader.IngestMode) { @@ -184,13 +185,7 @@ public class TradeHistoryWorker : BackgroundService } } - if (isWeeklyBiopsy) - { - // "NUR durch den TraderTraitCalculator schicken, NICHT persistieren." - // This would require resolving markets and positions and calling TraderTraitCalculator.Compute - // For now we skip persisting. - fetchedTrades = new List(); - } + // fetchedTrades are kept for weekly biopsy resolving, we will clear newTrades afterwards. } var newTrades = new List(); @@ -293,6 +288,26 @@ public class TradeHistoryWorker : BackgroundService newTrades = aggregated; } + if (isWeeklyBiopsy && newTrades.Count > 0) + { + var db = scope.ServiceProvider.GetRequiredService(); + var positions = await db.TraderPositions + .Include(p => p.MarketOutcome).ThenInclude(o => o.Market) + .Where(p => p.TraderId == trader.Id) + .ToListAsync(ct); + + var computedTraits = Predictalytics.Application.Services.TraderTraitCalculator.Compute(trader, newTrades, positions); + + db.TraderTraits.RemoveRange(db.TraderTraits.Where(tt => tt.TraderId == trader.Id)); + foreach (var (traitName, value) in computedTraits) + { + db.TraderTraits.Add(new Predictalytics.Domain.Entities.TraderTrait { TraderId = trader.Id, Trait = traitName, Value = value }); + } + await db.SaveChangesAsync(ct); + + newTrades.Clear(); // DO NOT persist trades! + } + if (trader.IngestMode == IngestMode.SnapshotOnly) { // "SnapshotOnly (Tier C): Stündlich: GetTraderPositionsAsync -> TraderPositions upserten" @@ -301,10 +316,12 @@ public class TradeHistoryWorker : BackgroundService try { var positions = await provider.GetTraderPositionsAsync(trader.PlatformUserId, ct); + decimal totalRealizedPnl = 0; + decimal totalUnrealizedPnl = 0; + if (positions != null && positions.Count > 0) { var db = scope.ServiceProvider.GetRequiredService(); - // For simplicity, just use the endpoint's positions var existingPos = db.TraderPositions.Where(tp => tp.TraderId == trader.Id).ToList(); foreach(var info in positions) { @@ -325,26 +342,123 @@ public class TradeHistoryWorker : BackgroundService { ex.SharesHeld = info.Size; ex.AvgCost = info.AveragePrice; - // RealizedPnl is built by tape replay, we skip it for SnapshotOnly + ex.RealizedPnl = info.RealizedPnl; // Mapped cashPnl! } else { - db.TraderPositions.Add(new Predictalytics.Domain.Entities.TraderPosition + ex = new Predictalytics.Domain.Entities.TraderPosition { TraderId = trader.Id, MarketOutcomeId = marketOutcomeId.Value, SharesHeld = info.Size, AvgCost = info.AveragePrice, - RealizedPnl = 0 - }); + RealizedPnl = info.RealizedPnl + }; + db.TraderPositions.Add(ex); } + + // Calculate unrealized PnL: pos.SharesHeld * (currentPrice - pos.AvgCost) + if (!string.IsNullOrEmpty(info.AssetId)) + { + var outcomeObj = await marketRepo.GetOutcomeByTokenIdAsync(info.AssetId, ct); + if (outcomeObj != null) + { + var unrealized = ex.SharesHeld * (outcomeObj.CurrentPrice - ex.AvgCost); + totalUnrealizedPnl += unrealized; + } + } + totalRealizedPnl += ex.RealizedPnl; } await db.SaveChangesAsync(ct); } + + // Overall PnL + decimal overallPnl = totalRealizedPnl + totalUnrealizedPnl; + + // Let's query Leaderboard PnLs + decimal pnl30d = 0; + decimal pnl7d = 0; + decimal pnl24h = 0; + + var polyApi = scope.ServiceProvider.GetService(); + if (polyApi != null) + { + try + { + var leaderboardAll = await polyApi.GetLeaderboardAsync(limit: 50, timePeriod: "ALL", ct: ct); + var entryAll = leaderboardAll.FirstOrDefault(e => string.Equals(e.ProxyWallet, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase) || string.Equals(e.UserName, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase)); + if (entryAll != null) + { + overallPnl = (decimal)entryAll.Pnl; + } + + var leaderboard30d = await polyApi.GetLeaderboardAsync(limit: 50, timePeriod: "30D", ct: ct); + var entry30d = leaderboard30d.FirstOrDefault(e => string.Equals(e.ProxyWallet, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase) || string.Equals(e.UserName, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase)); + if (entry30d != null) + { + pnl30d = (decimal)entry30d.Pnl; + } + + var leaderboard7d = await polyApi.GetLeaderboardAsync(limit: 50, timePeriod: "7D", ct: ct); + var entry7d = leaderboard7d.FirstOrDefault(e => string.Equals(e.ProxyWallet, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase) || string.Equals(e.UserName, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase)); + if (entry7d != null) + { + pnl7d = (decimal)entry7d.Pnl; + } + + var leaderboard24h = await polyApi.GetLeaderboardAsync(limit: 50, timePeriod: "24H", ct: ct); + var entry24h = leaderboard24h.FirstOrDefault(e => string.Equals(e.ProxyWallet, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase) || string.Equals(e.UserName, trader.PlatformUserId, StringComparison.OrdinalIgnoreCase)); + if (entry24h != null) + { + pnl24h = (decimal)entry24h.Pnl; + } + } + catch (Exception ex) + { + _logger.LogWarning(ex, "Failed to fetch leaderboard PnLs for SnapshotOnly trader {TraderId}", trader.Id); + } + } + + var dbCtx = scope.ServiceProvider.GetRequiredService(); + var analyticsObj = await dbCtx.TraderAnalytics.FirstOrDefaultAsync(a => a.TraderId == trader.Id, ct); + if (analyticsObj == null) + { + analyticsObj = new Predictalytics.Domain.Entities.TraderAnalytics { TraderId = trader.Id }; + dbCtx.TraderAnalytics.Add(analyticsObj); + } + + analyticsObj.OverallPnL = overallPnl; + analyticsObj.PnL30d = pnl30d; + analyticsObj.PnL7d = pnl7d; + analyticsObj.PnL24h = pnl24h; + analyticsObj.LastCalculatedAt = DateTime.UtcNow; + + trader.TotalPnl = overallPnl; + trader.LastAnalyzedAt = DateTime.UtcNow; + + // Save Daily Snapshot (Equity curve) + var today = DateTime.UtcNow.Date; + var snapshot = await dbCtx.TraderDailySnapshots.FirstOrDefaultAsync(s => s.TraderId == trader.Id && s.Date == today, ct); + if (snapshot == null) + { + dbCtx.TraderDailySnapshots.Add(new Predictalytics.Domain.Entities.TraderDailySnapshot + { + TraderId = trader.Id, + Date = today, + TotalPnl = overallPnl, + CurrentBalance = analyticsObj.CurrentBalance + }); + } + else + { + snapshot.TotalPnl = overallPnl; + } + + await dbCtx.SaveChangesAsync(ct); } catch (Exception ex) { - _logger.LogWarning(ex, "Failed to fetch positions for SnapshotOnly trader {TraderId}", trader.Id); + _logger.LogWarning(ex, "Failed to fetch positions / leaderboard for SnapshotOnly trader {TraderId}", trader.Id); } } }