using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; using FinlyticAssets.Models; using FinlyticAssets.Util; using FinlyticCore.Dtos.TradeRepublic; using FinlyticCore.Models.Assets; using FinlyticCore.Services; using FinlyticCore.Services.TradeRepublic; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; namespace FinlyticAssets.Services; /// /// A background service that runs a continuous asset synchronization loop, /// scanning Trade Republic to retrieve, update, and index all supported asset types. /// public class AssetScannerBackgroundService : BackgroundService { private readonly IServiceScopeFactory _serviceScopeFactory; private readonly IFinlyticLogger _finlyticLogger; private AssetsCount? _assetsCount; private AssetsCount? _currAssetsCount; public AssetScannerBackgroundService(IServiceScopeFactory serviceScopeFactory, IFinlyticLogger finlyticLogger) { _serviceScopeFactory = serviceScopeFactory; _finlyticLogger = finlyticLogger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] AssetScannerBackgroundService has started."); try { using var scope = _serviceScopeFactory.CreateScope(); var indexService = scope.ServiceProvider.GetRequiredService(); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Building initial asset index on service startup..."); await indexService.ReCreateIndexFileAsync(stoppingToken); } catch (Exception ex) { await _finlyticLogger.LogErrorAsync(SettingKeys.AssetsChannel, ex, "[AssetScannerBackgroundService] Failed to build initial asset index on startup. Continuing service execution."); } do { try { using var scope = _serviceScopeFactory.CreateScope(); var tradeRepublicService = scope.ServiceProvider.GetRequiredService(); var settingsService = scope.ServiceProvider.GetRequiredService(); var assetsDbService = scope.ServiceProvider.GetRequiredService(); var indexService = scope.ServiceProvider.GetRequiredService(); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Requesting total asset counts from Trade Republic..."); _assetsCount = await tradeRepublicService.GetAssetsCount(stoppingToken); _currAssetsCount = new AssetsCount(); var initSettings = await settingsService.GetSettings(); var isRecoveryMode = initSettings.CurrentScanningPage > 0; foreach (var type in Enum.GetValues()) { if (stoppingToken.IsCancellationRequested) break; if (isRecoveryMode) { if (type != initSettings.CurrentScanningType) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Recovery: {AssetType} was already processed. Skipping.", type); continue; } isRecoveryMode = false; } else { var settings = await settingsService.GetSettings(); settings.CurrentScanningType = type; settings.CurrentScanningPage = 0; await settingsService.SaveSettings(settings); } await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Processing asset type: {AssetType}...", type); await HandleAssetType(type, tradeRepublicService, settingsService, assetsDbService, indexService, stoppingToken); var currentSettings = await settingsService.GetSettings(); var delaySeconds = currentSettings.FinishedInitialScan ? currentSettings.AssetUpdateTypeDelay : currentSettings.InitAssetUpdateTypeDelay; if (delaySeconds > 0) { var jitter = Random.Shared.Next(0, Math.Min(15, delaySeconds)); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Waiting {Delay}s before next asset type ({Type}).", delaySeconds + jitter, type); await Task.Delay(TimeSpan.FromSeconds(delaySeconds + jitter), stoppingToken); } } var finalSettings = await settingsService.GetSettings(); finalSettings.CurrentScanningPage = 0; if (!finalSettings.FinishedInitialScan && !stoppingToken.IsCancellationRequested) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Initial scan successfully completed. Switching FinishedInitialScan to true."); finalSettings.FinishedInitialScan = true; } await settingsService.SaveSettings(finalSettings); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Full scan cycle completed. Waiting 1 minute before starting the next cycle."); await Task.Delay(TimeSpan.FromMinutes(1), stoppingToken); } catch (Exception e) when (!stoppingToken.IsCancellationRequested) { await _finlyticLogger.LogErrorAsync(SettingKeys.AssetsChannel, e, "[AssetScannerBackgroundService] An unhandled exception occurred in AssetScannerBackgroundService. Retrying in 10 seconds."); try { await Task.Delay(TimeSpan.FromSeconds(10), stoppingToken); } catch { /* Ignore */ } } } while (!stoppingToken.IsCancellationRequested); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] AssetScannerBackgroundService is stopping."); } private async Task HandleAssetType( AssetType type, ITradeRepublicService tradeRepublicService, ISettingsDbService settingsDbService, IAssetsDbService assetsDbService, IAssetsIndexService indexService, CancellationToken stoppingToken) { var totalCount = _assetsCount?.GetCountFromType(type) ?? 0; if (totalCount == 0) { await _finlyticLogger.LogWarningAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] No assets found for type {AssetType}.", type); return; } _currAssetsCount ??= new AssetsCount(); var currentItemOffset = 0; var settings = await settingsDbService.GetSettings(); var pageSize = Math.Clamp(settings.TradeRepublicMaxRequestPageSize <= 0 ? 50 : settings.TradeRepublicMaxRequestPageSize, 1, 100); if (settings.CurrentScanningType == type && settings.CurrentScanningPage > 0) { currentItemOffset = (settings.CurrentScanningPage - 1) * pageSize; await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Resuming full scan for {AssetType} from Page {Page} (Calculated Offset: {Offset}).", type, settings.CurrentScanningPage, currentItemOffset); } while (currentItemOffset < totalCount && !stoppingToken.IsCancellationRequested) { var currentSettings = await settingsDbService.GetSettings(); pageSize = Math.Clamp(currentSettings.TradeRepublicMaxRequestPageSize <= 0 ? 50 : currentSettings.TradeRepublicMaxRequestPageSize, 1, 100); var currentPage = (currentItemOffset / pageSize) + 1; currentSettings.CurrentScanningType = type; currentSettings.CurrentScanningPage = currentPage; await settingsDbService.SaveSettings(currentSettings); await _finlyticLogger.LogDebugAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Fetching {AssetType} - Page {Page}. Numerical Offset: {Offset}/{Total}", type, currentPage, currentItemOffset, totalCount); var assets = await tradeRepublicService.GetAssets(type, currentPage, pageSize, stoppingToken); if (assets?.Results == null || assets.Results.Count == 0) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Fetch for {AssetType} (Page {Page}) returned no results. Reached end of available assets.", type, currentPage); currentSettings.CurrentScanningPage = 0; await settingsDbService.SaveSettings(currentSettings); break; } await ProcessAssets(assets.Results, assetsDbService, indexService, stoppingToken); _currAssetsCount.SetCountOfType(type, _currAssetsCount.GetCountFromType(type) + assets.Results.Count); currentItemOffset += assets.Results.Count; if (assets.Results.Count < pageSize) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Reached the last page for {AssetType}.", type); currentSettings.CurrentScanningPage = 0; await settingsDbService.SaveSettings(currentSettings); break; } var delaySeconds = currentSettings.FinishedInitialScan ? currentSettings.BatchAssetUpdateDelay : currentSettings.InitBatchAssetUpdateDelay; if (delaySeconds > 0) { var maxJitter = Math.Min(5, delaySeconds); var jitter = Random.Shared.Next(0, maxJitter + 1); await _finlyticLogger.LogDebugAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Waiting {Delay} seconds before the next batch.", delaySeconds + jitter); await Task.Delay(TimeSpan.FromSeconds(delaySeconds + jitter), stoppingToken); } } await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Finished scanning {AssetType}. Total scanned in this cycle: {Count}/{Total}", type, _currAssetsCount.GetCountFromType(type), totalCount); } private async Task ProcessAssets( IList assets, IAssetsDbService assetsDbService, IAssetsIndexService indexService, CancellationToken stoppingToken) { if (assets == null || assets.Count == 0) return; var changedRows = await assetsDbService.AddOrUpdateAssetsAsync(assets); await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] [Scan] {Count} assets passed to the DB service. {Changed} modifications/inserts executed.", assets.Count, changedRows); if (changedRows > 0) { await _finlyticLogger.LogInfoAsync(SettingKeys.AssetsChannel, "[AssetScannerBackgroundService] Database modifications detected. Recreating the asset index file..."); await indexService.ReCreateIndexFileAsync(stoppingToken); } } }