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;
///
/// 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.
///
public class MarketSyncWorker : BackgroundService
{
private static DateTime _lastDbError = DateTime.MinValue;
private readonly IServiceProvider _services;
private readonly IPlatformStatisticsService _statsService;
private readonly ILogger _logger;
public MarketSyncWorker(IServiceProvider services, ILogger 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();
var platformProviders = _services.GetRequiredService>();
foreach (var p in platformProviders.Where(x => x.IsImplemented))
{
try
{
if (stoppingToken.IsCancellationRequested) break;
using var platformCtx = PlatformLogContext.Push(p.PlatformName);
int cycleTotalSynced = 0;
foreach (var includeClosed in new[] { false, true })
{
_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();
var provider = scope.ServiceProvider.GetRequiredService>()
.First(x => x.Platform == p.Platform);
await rateLimiter.WaitAsync(provider.Platform, stoppingToken);
var markets = await provider.GetMarketsAsync(batchSize, offset.ToString(), includeClosed, stoppingToken);
if (markets.Count == 0) break;
await marketRepo.AddOrUpdateRangeAsync(markets, stoppingToken);
passSynced += markets.Count;
cycleTotalSynced += markets.Count;
_statsService.TrackMarketSync(provider.Platform, markets.Count);
offset += batchSize;
if (passSynced % 500 == 0)
_logger.LogWarning("[{Platform}] Synced {Total} markets so far (includeClosed={Closed})...", p.PlatformName, passSynced, includeClosed);
}
}
// Need a temporary scope for stats
using (var scope = _services.CreateScope())
{
var marketRepo = scope.ServiceProvider.GetRequiredService();
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");
}
}