diff --git a/src/Predictalytics.Domain/Interfaces/ITradeRepository.cs b/src/Predictalytics.Domain/Interfaces/ITradeRepository.cs index 3500cd4..b544427 100644 --- a/src/Predictalytics.Domain/Interfaces/ITradeRepository.cs +++ b/src/Predictalytics.Domain/Interfaces/ITradeRepository.cs @@ -17,6 +17,6 @@ public interface ITradeRepository Task AddRangeAsync(IEnumerable trades, CancellationToken ct = default); Task GetTotalVolumeAsync(DateTime? since = null, CancellationToken ct = default); Task> GetOrphanedTradesAsync(int limit, CancellationToken ct = default); - Task> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, CancellationToken ct = default); + Task> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, IEnumerable platformTradeIds, CancellationToken ct = default); Task UpdateAsync(Trade trade, CancellationToken ct = default); } diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/MarketRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/MarketRepository.cs index 898ca23..6c0f47e 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/MarketRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/MarketRepository.cs @@ -28,21 +28,29 @@ public class MarketRepository : IMarketRepository public async Task AddOrUpdateAsync(Market market, CancellationToken ct = default) { - TruncateMarketStrings(market); - - var existing = await _db.Markets.Include(m => m.Outcomes) - .FirstOrDefaultAsync(m => m.Platform == market.Platform && m.PlatformMarketId == market.PlatformMarketId, ct); - - if (existing != null) + await _syncSemaphore.WaitAsync(ct); + try { - UpdateMarketFields(existing, market); - } - else - { - _db.Markets.Add(market); - } + TruncateMarketStrings(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 markets, CancellationToken ct = default) diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs index 129ddea..71f0fed 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs @@ -89,10 +89,13 @@ public class TradeRepository : ITradeRepository .ToListAsync(ct); } - public async Task> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, CancellationToken ct = default) + public async Task> GetKnownPlatformTradeIdsAsync(PlatformType platform, int traderId, IEnumerable platformTradeIds, CancellationToken ct = default) { + var idList = platformTradeIds.ToList(); + if (idList.Count == 0) return new HashSet(); + 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) .ToListAsync(ct); return new HashSet(ids); diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs index aeb42d6..39df856 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/TraderRepository.cs @@ -92,6 +92,7 @@ public class TraderRepository : ITraderRepository public async Task> GetTradersForCleanupAsync(DateTime inactiveSince, DateTime errorSince, int take = 50, CancellationToken ct = default) { return await _db.Traders + .Where(t => t.IsAutoDiscovered && !t.WatchlistEntries.Any()) .Where(t => (t.LastPolledAt != null && t.LastPolledAt < inactiveSince) || (t.LastApiErrorAt != null && t.LastApiErrorAt < errorSince)) .OrderBy(t => t.LastApiErrorAt ?? DateTime.MaxValue) // Prioritize errors first diff --git a/src/Predictalytics.Worker/Services/MarketSyncWorker.cs b/src/Predictalytics.Worker/Services/MarketSyncWorker.cs index 0166586..04c299b 100644 --- a/src/Predictalytics.Worker/Services/MarketSyncWorker.cs +++ b/src/Predictalytics.Worker/Services/MarketSyncWorker.cs @@ -15,6 +15,7 @@ namespace Predictalytics.Worker.Services; public class MarketSyncWorker : BackgroundService { private static DateTime _lastDbError = DateTime.MinValue; + private static DateTime _lastClosedMarketSync = DateTime.MinValue; private readonly IServiceProvider _services; private readonly IPlatformStatisticsService _statsService; private readonly ILogger _logger; @@ -42,7 +43,13 @@ public class MarketSyncWorker : BackgroundService using var platformCtx = PlatformLogContext.Push(p.PlatformName); int cycleTotalSynced = 0; - foreach (var includeClosed in new[] { false, true }) + var includeClosedOptions = new List { 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); @@ -73,6 +80,11 @@ public class MarketSyncWorker : BackgroundService if (passSynced % 500 == 0) _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 diff --git a/src/Predictalytics.Worker/Services/PollingWorker.cs b/src/Predictalytics.Worker/Services/PollingWorker.cs index 386cb0e..a3f3346 100644 --- a/src/Predictalytics.Worker/Services/PollingWorker.cs +++ b/src/Predictalytics.Worker/Services/PollingWorker.cs @@ -77,7 +77,8 @@ public class PollingWorker : BackgroundService // ── Deduplicate: check each trade against DB ── var newTrades = new List(); - 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) { diff --git a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs index 4fcaf38..83b5d7d 100644 --- a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs @@ -85,7 +85,8 @@ public class TradeHistoryWorker : BackgroundService var fetchedTrades = await provider.GetTraderTradesAsync(trader.PlatformUserId, TradesPerFetch, ct); 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 var assetIdsToResolve = validTrades @@ -111,7 +112,6 @@ public class TradeHistoryWorker : BackgroundService { if (knownTradeIds.Contains(trade.PlatformTradeId)) { - if (!isInitial) break; continue; }