D3: Implement IngestMode classification, weekly biopsy, SnapshotOnly bypass, and Aggregated import grouping

This commit is contained in:
Richard
2026-07-19 17:23:19 +02:00
parent 6fd7e563f4
commit 3ed0b4df27
7 changed files with 401 additions and 17 deletions
@@ -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<AppDbContext>()
.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<ITraderRepository, TraderRepository>();
var mockTradeRepo = new Mock<ITradeRepository>();
mockTradeRepo.Setup(r => r.GetKnownPlatformTradeIdsAsync(It.IsAny<PlatformType>(), It.IsAny<int>(), It.IsAny<IEnumerable<string>>(), It.IsAny<CancellationToken>()))
.ReturnsAsync(new HashSet<string>());
services.AddSingleton(mockTradeRepo.Object);
services.AddScoped<IMarketRepository, MarketRepository>();
services.AddSingleton(typeof(ILogger<>), typeof(NullLogger<>));
var mockProvider = new Mock<IPlatformProvider>();
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<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()))
.ReturnsAsync(new List<Trade> { new Trade { PlatformTradeId = "test_tx" } });
services.AddSingleton<IEnumerable<IPlatformProvider>>(new[] { mockProvider.Object });
var mockRateLimiter = new Mock<IRateLimiter>();
services.AddSingleton(mockRateLimiter.Object);
var mockStats = new Mock<IPlatformStatisticsService>();
services.AddSingleton(mockStats.Object);
using var provider = services.BuildServiceProvider();
var worker = new PollingWorker(provider, NullLogger<PollingWorker>.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<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()), Times.Never);
}
[Fact]
public async Task AggregatedMode_AggregatesTradesCorrectly()
{
using var connection = new SqliteConnection("DataSource=:memory:");
connection.Open();
var options = new DbContextOptionsBuilder<AppDbContext>()
.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<Event>().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<ITraderRepository, TraderRepository>();
var mockTradeRepo = new Mock<ITradeRepository>();
mockTradeRepo.Setup(r => r.GetKnownPlatformTradeIdsAsync(It.IsAny<PlatformType>(), It.IsAny<int>(), It.IsAny<IEnumerable<string>>(), It.IsAny<CancellationToken>()))
.ReturnsAsync(new HashSet<string>());
// Mock AddRangeAsync to bypass MySql raw query and save trades directly to SQLite in-memory DB
mockTradeRepo.Setup(r => r.AddRangeAsync(It.IsAny<IEnumerable<Trade>>(), It.IsAny<CancellationToken>()))
.Callback<IEnumerable<Trade>, 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<IMarketRepository, MarketRepository>();
services.AddSingleton(typeof(ILogger<>), typeof(NullLogger<>));
var baseTime = new DateTime(2026, 7, 19, 10, 30, 0, DateTimeKind.Utc);
var mockProvider = new Mock<IPlatformProvider>();
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<string>(), It.IsAny<int>(), It.IsAny<CancellationToken>()))
.ReturnsAsync(new List<Trade>
{
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<IEnumerable<IPlatformProvider>>(new[] { mockProvider.Object });
var mockRateLimiter = new Mock<IRateLimiter>();
services.AddSingleton(mockRateLimiter.Object);
var mockStats = new Mock<IPlatformStatisticsService>();
services.AddSingleton(mockStats.Object);
using var provider = services.BuildServiceProvider();
var worker = new PollingWorker(provider, NullLogger<PollingWorker>.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);
}
}
}
@@ -11,6 +11,11 @@ public static class TraderTraitCalculator
IReadOnlyCollection<TraderPosition> positions) IReadOnlyCollection<TraderPosition> positions)
{ {
var traits = new List<(string, decimal)>(); var traits = new List<(string, decimal)>();
if (trader.IngestMode == IngestMode.SnapshotOnly)
{
traits.Add(("not_copyable_hf", 1.0m));
}
if (trades.Count == 0) return traits; if (trades.Count == 0) return traits;
var now = DateTime.UtcNow; var now = DateTime.UtcNow;
@@ -58,7 +58,8 @@ public record TraderPositionInfo(
decimal AveragePrice, decimal AveragePrice,
decimal CurrentValue, decimal CurrentValue,
decimal PnlPercent, decimal PnlPercent,
string? AssetId = null string? AssetId = null,
decimal RealizedPnl = 0
); );
/// <summary> /// <summary>
@@ -151,7 +151,7 @@ public class PolymarketProvider : IPlatformProvider
platformUserId, r.Market, r.Question, r.Outcome, platformUserId, r.Market, r.Question, r.Outcome,
(decimal)r.Size, (decimal)r.AvgPrice, (decimal)r.Size, (decimal)r.AvgPrice,
(decimal)r.CurrentValue, (decimal)r.PercentPnl, (decimal)r.CurrentValue, (decimal)r.PercentPnl,
r.AssetId r.AssetId, (decimal)r.CashPnl
)).ToList(); )).ToList();
} }
@@ -35,6 +35,12 @@ public class PositionPnLEngine : IPositionPnLEngine
return; 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 // Fetch all trades for this trader, sorted chronologically, including outcomes and markets
var trades = await _db.Trades var trades = await _db.Trades
.Include(t => t.MarketOutcome) .Include(t => t.MarketOutcome)
@@ -53,7 +53,7 @@ public class PollingWorker : BackgroundService
var rateLimiter = scope.ServiceProvider.GetRequiredService<IRateLimiter>(); var rateLimiter = scope.ServiceProvider.GetRequiredService<IRateLimiter>();
var trader = await traderRepo.GetByIdAsync(t.Id, stoppingToken); 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); var provider = providers.FirstOrDefault(p => p.Platform == trader.Platform && p.IsImplemented);
if (provider == null) continue; if (provider == null) continue;
@@ -122,6 +122,37 @@ public class PollingWorker : BackgroundService
} }
} }
if (trader.IngestMode == IngestMode.Aggregated && newTrades.Count > 0)
{
var aggregated = new List<Domain.Entities.Trade>();
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 ── // ── Persist new trades ──
if (newTrades.Count > 0) if (newTrades.Count > 0)
{ {
@@ -6,6 +6,7 @@ using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Configuration;
using Microsoft.EntityFrameworkCore;
namespace Predictalytics.Worker.Services; namespace Predictalytics.Worker.Services;
@@ -173,7 +174,7 @@ public class TradeHistoryWorker : BackgroundService
if (tradesPerDay > 5000) newMode = IngestMode.SnapshotOnly; if (tradesPerDay > 5000) newMode = IngestMode.SnapshotOnly;
else if (tradesPerDay > 100 && trader.IngestMode == IngestMode.Full) newMode = IngestMode.Aggregated; 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.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) if (newMode != trader.IngestMode)
{ {
@@ -184,13 +185,7 @@ public class TradeHistoryWorker : BackgroundService
} }
} }
if (isWeeklyBiopsy) // fetchedTrades are kept for weekly biopsy resolving, we will clear newTrades afterwards.
{
// "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<Domain.Entities.Trade>();
}
} }
var newTrades = new List<Domain.Entities.Trade>(); var newTrades = new List<Domain.Entities.Trade>();
@@ -293,6 +288,26 @@ public class TradeHistoryWorker : BackgroundService
newTrades = aggregated; newTrades = aggregated;
} }
if (isWeeklyBiopsy && newTrades.Count > 0)
{
var db = scope.ServiceProvider.GetRequiredService<Predictalytics.Infrastructure.Data.AppDbContext>();
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) if (trader.IngestMode == IngestMode.SnapshotOnly)
{ {
// "SnapshotOnly (Tier C): Stündlich: GetTraderPositionsAsync -> TraderPositions upserten" // "SnapshotOnly (Tier C): Stündlich: GetTraderPositionsAsync -> TraderPositions upserten"
@@ -301,10 +316,12 @@ public class TradeHistoryWorker : BackgroundService
try try
{ {
var positions = await provider.GetTraderPositionsAsync(trader.PlatformUserId, ct); var positions = await provider.GetTraderPositionsAsync(trader.PlatformUserId, ct);
decimal totalRealizedPnl = 0;
decimal totalUnrealizedPnl = 0;
if (positions != null && positions.Count > 0) if (positions != null && positions.Count > 0)
{ {
var db = scope.ServiceProvider.GetRequiredService<Predictalytics.Infrastructure.Data.AppDbContext>(); var db = scope.ServiceProvider.GetRequiredService<Predictalytics.Infrastructure.Data.AppDbContext>();
// For simplicity, just use the endpoint's positions
var existingPos = db.TraderPositions.Where(tp => tp.TraderId == trader.Id).ToList(); var existingPos = db.TraderPositions.Where(tp => tp.TraderId == trader.Id).ToList();
foreach(var info in positions) foreach(var info in positions)
{ {
@@ -325,26 +342,123 @@ public class TradeHistoryWorker : BackgroundService
{ {
ex.SharesHeld = info.Size; ex.SharesHeld = info.Size;
ex.AvgCost = info.AveragePrice; ex.AvgCost = info.AveragePrice;
// RealizedPnl is built by tape replay, we skip it for SnapshotOnly ex.RealizedPnl = info.RealizedPnl; // Mapped cashPnl!
} }
else else
{ {
db.TraderPositions.Add(new Predictalytics.Domain.Entities.TraderPosition ex = new Predictalytics.Domain.Entities.TraderPosition
{ {
TraderId = trader.Id, TraderId = trader.Id,
MarketOutcomeId = marketOutcomeId.Value, MarketOutcomeId = marketOutcomeId.Value,
SharesHeld = info.Size, SharesHeld = info.Size,
AvgCost = info.AveragePrice, 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); 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<Predictalytics.Infrastructure.Providers.Polymarket.PolymarketApiClient>();
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) catch (Exception ex)
{ {
_logger.LogWarning(ex, "Failed to fetch positions for SnapshotOnly trader {TraderId}", trader.Id); _logger.LogWarning(ex, "Failed to fetch leaderboard PnLs for SnapshotOnly trader {TraderId}", trader.Id);
}
}
var dbCtx = scope.ServiceProvider.GetRequiredService<Predictalytics.Infrastructure.Data.AppDbContext>();
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 / leaderboard for SnapshotOnly trader {TraderId}", trader.Id);
} }
} }
} }