Files
Predictalytics/src/Predictalytics.Worker/Services/MarketSyncWorker.cs
T

124 lines
5.9 KiB
C#

using Predictalytics.Application.Interfaces;
using Predictalytics.Domain.Enums;
using Predictalytics.Domain.Interfaces;
using Predictalytics.Infrastructure.Logging;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace Predictalytics.Worker.Services;
/// <summary>
/// Background service that periodically syncs market master data (including outcomes/token IDs)
/// from the Gamma API. This builds the lookup table needed to resolve trades to markets.
/// </summary>
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<MarketSyncWorker> _logger;
public MarketSyncWorker(IServiceProvider services, ILogger<MarketSyncWorker> logger, IPlatformStatisticsService statsService)
{ _services = services; _logger = logger; _statsService = statsService; }
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("🏛️ MarketSyncWorker started");
await Task.Delay(5000, stoppingToken); // Let DB initialize
while (!stoppingToken.IsCancellationRequested)
{
try
{
var rateLimiter = _services.GetRequiredService<IRateLimiter>();
var platformProviders = _services.GetRequiredService<IEnumerable<IPlatformProvider>>();
foreach (var p in platformProviders.Where(x => x.IsImplemented))
{
try
{
if (stoppingToken.IsCancellationRequested) break;
using var platformCtx = PlatformLogContext.Push(p.PlatformName);
int cycleTotalSynced = 0;
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);
int passSynced = 0;
int offset = 0;
const int batchSize = 1000;
while (!stoppingToken.IsCancellationRequested)
{
// Fresh scope per batch
using var scope = _services.CreateScope();
var marketRepo = scope.ServiceProvider.GetRequiredService<IMarketRepository>();
var provider = scope.ServiceProvider.GetRequiredService<IEnumerable<IPlatformProvider>>()
.First(x => x.Platform == p.Platform);
await rateLimiter.WaitAsync(provider.Platform, stoppingToken);
var events = await provider.GetEventsAsync(batchSize, offset.ToString(), includeClosed, stoppingToken);
if (events.Count == 0) break;
await marketRepo.AddOrUpdateEventsAsync(events, stoppingToken);
passSynced += events.Count;
cycleTotalSynced += events.Count;
_statsService.TrackMarketSync(provider.Platform, events.Count);
offset += batchSize;
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
using (var scope = _services.CreateScope())
{
var marketRepo = scope.ServiceProvider.GetRequiredService<IMarketRepository>();
var dbCount = await marketRepo.GetCountAsync(stoppingToken);
_logger.LogWarning("✅ [{Platform}] Market sync complete. {Synced} synced this cycle, {Total} total in DB",
p.PlatformName, cycleTotalSynced, dbCount);
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error syncing markets for platform {Platform}", p.PlatformName);
}
}
}
catch (OperationCanceledException) { break; }
catch (Exception ex) when (ex.ToString().Contains("MySqlException") || ex.ToString().Contains("Connection"))
{
if (DateTime.UtcNow - _lastDbError > TimeSpan.FromMinutes(10))
{
_logger.LogWarning("⚠️ Database connection lost in MarketSyncWorker. Retrying in 30m. (Error: {Message})", ex.Message);
_lastDbError = DateTime.UtcNow;
}
}
catch (Exception ex) { _logger.LogError(ex, "MarketSyncWorker error"); }
// Run every 30 minutes
_logger.LogInformation("🏛️ Next market sync in 30 minutes.");
await Task.Delay(TimeSpan.FromMinutes(30), stoppingToken);
}
_logger.LogInformation("🏛️ MarketSyncWorker stopped");
}
}