Fix B6: Market concurrency lock, watchlist exclusion from cleanup, optimized trade ID deduplication query, and TradeHistoryWorker sync logic

This commit is contained in:
Richard
2026-07-03 11:10:53 +02:00
parent 434c631290
commit 5fe7ab63ac
7 changed files with 45 additions and 20 deletions
@@ -17,6 +17,6 @@ public interface ITradeRepository
Task AddRangeAsync(IEnumerable<Trade> trades, CancellationToken ct = default); Task AddRangeAsync(IEnumerable<Trade> trades, CancellationToken ct = default);
Task<decimal> GetTotalVolumeAsync(DateTime? since = null, CancellationToken ct = default); Task<decimal> GetTotalVolumeAsync(DateTime? since = null, CancellationToken ct = default);
Task<IReadOnlyList<Trade>> GetOrphanedTradesAsync(int limit, CancellationToken ct = default); Task<IReadOnlyList<Trade>> GetOrphanedTradesAsync(int limit, CancellationToken ct = default);
Task<HashSet<string>> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, CancellationToken ct = default); Task<HashSet<string>> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, IEnumerable<string> platformTradeIds, CancellationToken ct = default);
Task UpdateAsync(Trade trade, CancellationToken ct = default); Task UpdateAsync(Trade trade, CancellationToken ct = default);
} }
@@ -28,21 +28,29 @@ public class MarketRepository : IMarketRepository
public async Task AddOrUpdateAsync(Market market, CancellationToken ct = default) public async Task AddOrUpdateAsync(Market market, CancellationToken ct = default)
{ {
TruncateMarketStrings(market); await _syncSemaphore.WaitAsync(ct);
try
var existing = await _db.Markets.Include(m => m.Outcomes)
.FirstOrDefaultAsync(m => m.Platform == market.Platform && m.PlatformMarketId == market.PlatformMarketId, ct);
if (existing != null)
{ {
UpdateMarketFields(existing, market); TruncateMarketStrings(market);
}
else
{
_db.Markets.Add(market);
}
await _db.SaveChangesAsync(ct); var existing = await _db.Markets.Include(m => m.Outcomes)
.FirstOrDefaultAsync(m => m.Platform == market.Platform && m.PlatformMarketId == market.PlatformMarketId, ct);
if (existing != null)
{
UpdateMarketFields(existing, market);
}
else
{
_db.Markets.Add(market);
}
await _db.SaveChangesAsync(ct);
}
finally
{
_syncSemaphore.Release();
}
} }
public async Task AddOrUpdateRangeAsync(IEnumerable<Market> markets, CancellationToken ct = default) public async Task AddOrUpdateRangeAsync(IEnumerable<Market> markets, CancellationToken ct = default)
@@ -89,10 +89,13 @@ public class TradeRepository : ITradeRepository
.ToListAsync(ct); .ToListAsync(ct);
} }
public async Task<HashSet<string>> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, CancellationToken ct = default) public async Task<HashSet<string>> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, IEnumerable<string> platformTradeIds, CancellationToken ct = default)
{ {
var idList = platformTradeIds.ToList();
if (idList.Count == 0) return new HashSet<string>();
var ids = await _db.Trades var ids = await _db.Trades
.Where(t => t.Platform == platform && t.TraderId == traderId) .Where(t => t.Platform == platform && t.TraderId == traderId && idList.Contains(t.PlatformTradeId))
.Select(t => t.PlatformTradeId) .Select(t => t.PlatformTradeId)
.ToListAsync(ct); .ToListAsync(ct);
return new HashSet<string>(ids); return new HashSet<string>(ids);
@@ -92,6 +92,7 @@ public class TraderRepository : ITraderRepository
public async Task<IReadOnlyList<Trader>> GetTradersForCleanupAsync(DateTime inactiveSince, DateTime errorSince, int take = 50, CancellationToken ct = default) public async Task<IReadOnlyList<Trader>> GetTradersForCleanupAsync(DateTime inactiveSince, DateTime errorSince, int take = 50, CancellationToken ct = default)
{ {
return await _db.Traders return await _db.Traders
.Where(t => t.IsAutoDiscovered && !t.WatchlistEntries.Any())
.Where(t => (t.LastPolledAt != null && t.LastPolledAt < inactiveSince) || .Where(t => (t.LastPolledAt != null && t.LastPolledAt < inactiveSince) ||
(t.LastApiErrorAt != null && t.LastApiErrorAt < errorSince)) (t.LastApiErrorAt != null && t.LastApiErrorAt < errorSince))
.OrderBy(t => t.LastApiErrorAt ?? DateTime.MaxValue) // Prioritize errors first .OrderBy(t => t.LastApiErrorAt ?? DateTime.MaxValue) // Prioritize errors first
@@ -15,6 +15,7 @@ namespace Predictalytics.Worker.Services;
public class MarketSyncWorker : BackgroundService public class MarketSyncWorker : BackgroundService
{ {
private static DateTime _lastDbError = DateTime.MinValue; private static DateTime _lastDbError = DateTime.MinValue;
private static DateTime _lastClosedMarketSync = DateTime.MinValue;
private readonly IServiceProvider _services; private readonly IServiceProvider _services;
private readonly IPlatformStatisticsService _statsService; private readonly IPlatformStatisticsService _statsService;
private readonly ILogger<MarketSyncWorker> _logger; private readonly ILogger<MarketSyncWorker> _logger;
@@ -42,7 +43,13 @@ public class MarketSyncWorker : BackgroundService
using var platformCtx = PlatformLogContext.Push(p.PlatformName); using var platformCtx = PlatformLogContext.Push(p.PlatformName);
int cycleTotalSynced = 0; int cycleTotalSynced = 0;
foreach (var includeClosed in new[] { false, true }) var includeClosedOptions = new List<bool> { false };
if (DateTime.UtcNow - _lastClosedMarketSync >= TimeSpan.FromDays(1))
{
includeClosedOptions.Add(true);
}
foreach (var includeClosed in includeClosedOptions)
{ {
_logger.LogWarning("[{Platform}] Syncing markets (includeClosed={Closed})...", p.PlatformName, includeClosed); _logger.LogWarning("[{Platform}] Syncing markets (includeClosed={Closed})...", p.PlatformName, includeClosed);
@@ -73,6 +80,11 @@ public class MarketSyncWorker : BackgroundService
if (passSynced % 500 == 0) if (passSynced % 500 == 0)
_logger.LogWarning("[{Platform}] Synced {Total} markets so far (includeClosed={Closed})...", p.PlatformName, passSynced, includeClosed); _logger.LogWarning("[{Platform}] Synced {Total} markets so far (includeClosed={Closed})...", p.PlatformName, passSynced, includeClosed);
} }
if (includeClosed)
{
_lastClosedMarketSync = DateTime.UtcNow;
}
} }
// Need a temporary scope for stats // Need a temporary scope for stats
@@ -77,7 +77,8 @@ public class PollingWorker : BackgroundService
// ── Deduplicate: check each trade against DB ── // ── Deduplicate: check each trade against DB ──
var newTrades = new List<Domain.Entities.Trade>(); var newTrades = new List<Domain.Entities.Trade>();
var knownTradeIds = await tradeRepo.GetKnownPlatformTradeIdsAsync(trader.Platform, trader.Id, stoppingToken); var fetchedTradeIds = validTrades.Select(tr => tr.PlatformTradeId).ToList();
var knownTradeIds = await tradeRepo.GetKnownPlatformTradeIdsAsync(trader.Platform, trader.Id, fetchedTradeIds, stoppingToken);
foreach (var trade in validTrades) foreach (var trade in validTrades)
{ {
@@ -85,7 +85,8 @@ public class TradeHistoryWorker : BackgroundService
var fetchedTrades = await provider.GetTraderTradesAsync(trader.PlatformUserId, TradesPerFetch, ct); var fetchedTrades = await provider.GetTraderTradesAsync(trader.PlatformUserId, TradesPerFetch, ct);
var validTrades = fetchedTrades.Where(tr => !string.IsNullOrWhiteSpace(tr.PlatformTradeId)).ToList(); var validTrades = fetchedTrades.Where(tr => !string.IsNullOrWhiteSpace(tr.PlatformTradeId)).ToList();
var knownTradeIds = await tradeRepo.GetKnownPlatformTradeIdsAsync(trader.Platform, trader.Id, ct); var fetchedTradeIds = validTrades.Select(tr => tr.PlatformTradeId).ToList();
var knownTradeIds = await tradeRepo.GetKnownPlatformTradeIdsAsync(trader.Platform, trader.Id, fetchedTradeIds, ct);
// Collect all unique AssetIds we might need to resolve // Collect all unique AssetIds we might need to resolve
var assetIdsToResolve = validTrades var assetIdsToResolve = validTrades
@@ -111,7 +112,6 @@ public class TradeHistoryWorker : BackgroundService
{ {
if (knownTradeIds.Contains(trade.PlatformTradeId)) if (knownTradeIds.Contains(trade.PlatformTradeId))
{ {
if (!isInitial) break;
continue; continue;
} }