using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using FinlyticCore.Dtos.TechnicalAnalysis; using FinlyticCore.Services.TradeRepublic; using FinlyticTechnicalAnalysis.Database; using FinlyticTechnicalAnalysis.Entities; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace FinlyticTechnicalAnalysis.Services; public interface ITechnicalAnalysisDbService { Task GetAnalysisAsync(string isin, bool forceRefresh = false, string? ticker = null, CancellationToken cancellationToken = default); Task GetLivePriceAsync(string isin, CancellationToken cancellationToken = default); } public class TechnicalAnalysisDbService : ITechnicalAnalysisDbService { private readonly IServiceScopeFactory _scopeFactory; private readonly IYahooMarketDataScraper _yahooScraper; private readonly ITradeRepublicService _trService; private readonly ITechnicalAnalysisCalculator _calculator; private readonly ILogger _logger; private static readonly ConcurrentDictionary Candles, string Symbol, string Currency, DateTime FetchedAt)> _candleCache = new(); private static readonly ConcurrentDictionary _perIsinLocks = new(); private static readonly TimeSpan CandleCacheTtl = TimeSpan.FromMinutes(15); private static readonly TimeSpan DbCacheTtl = TimeSpan.FromHours(1); public TechnicalAnalysisDbService( IServiceScopeFactory scopeFactory, IYahooMarketDataScraper yahooScraper, ITradeRepublicService trService, ITechnicalAnalysisCalculator calculator, ILogger logger) { _scopeFactory = scopeFactory; _yahooScraper = yahooScraper; _trService = trService; _calculator = calculator; _logger = logger; } public async Task GetAnalysisAsync(string isin, bool forceRefresh = false, string? ticker = null, CancellationToken cancellationToken = default) { if (string.IsNullOrWhiteSpace(isin)) return null; var cleanIsin = isin.Trim().ToUpperInvariant(); // 1. Layer-1: Fast-Path aus In-Memory Cache (wenn kein forceRefresh) if (!forceRefresh && _candleCache.TryGetValue(cleanIsin, out var ramEntry) && DateTime.UtcNow - ramEntry.FetchedAt < CandleCacheTtl && (string.IsNullOrWhiteSpace(ticker) || string.Equals(ramEntry.Symbol, ticker, StringComparison.OrdinalIgnoreCase))) { _logger.LogDebug("[{Channel}] RAM-Cache Hit for ISIN {Isin}. Merging live price...", "TechnicalAnalysisChannel", cleanIsin); return await BuildAnalysisWithLivePriceAsync(cleanIsin, ramEntry.Candles, ramEntry.Symbol, ramEntry.Currency, cancellationToken); } var semaphore = _perIsinLocks.GetOrAdd(cleanIsin, _ => new SemaphoreSlim(1, 1)); await semaphore.WaitAsync(cancellationToken); try { // Re-Check nach Lock-Erhalt if (!forceRefresh && _candleCache.TryGetValue(cleanIsin, out ramEntry) && DateTime.UtcNow - ramEntry.FetchedAt < CandleCacheTtl && (string.IsNullOrWhiteSpace(ticker) || string.Equals(ramEntry.Symbol, ticker, StringComparison.OrdinalIgnoreCase))) { return await BuildAnalysisWithLivePriceAsync(cleanIsin, ramEntry.Candles, ramEntry.Symbol, ramEntry.Currency, cancellationToken); } // 2. Layer-2: Prüfen ob frische Daten in der Datenbank liegen if (!forceRefresh) { var dbDto = await GetFromDbCacheAsync(cleanIsin, ticker, cancellationToken); if (dbDto != null) { _logger.LogDebug("[{Channel}] DB-Cache Hit for ISIN {Isin}.", "TechnicalAnalysisChannel", cleanIsin); return dbDto; } } return await FullRefreshAsync(cleanIsin, ticker, cancellationToken); } finally { semaphore.Release(); if (semaphore.CurrentCount == 1) { _perIsinLocks.TryRemove(cleanIsin, out _); } } } public async Task GetLivePriceAsync(string isin, CancellationToken cancellationToken = default) { if (string.IsNullOrWhiteSpace(isin)) return null; var cleanIsin = isin.Trim().ToUpperInvariant(); var (livePrice, liveBid, liveAsk, preChange) = await FetchLivePriceAsync(cleanIsin, cancellationToken); if (!livePrice.HasValue) return null; return new LivePriceDto( cleanIsin, Math.Round(livePrice.Value, 2), preChange ?? 0m, liveBid.HasValue ? Math.Round(liveBid.Value, 2) : null, liveAsk.HasValue ? Math.Round(liveAsk.Value, 2) : null ); } private async Task FullRefreshAsync(string cleanIsin, string? requestedTicker, CancellationToken cancellationToken) { _logger.LogInformation("[{Channel}] Full refresh for ISIN {Isin} (RequestedTicker: {Ticker})", "TechnicalAnalysisChannel", cleanIsin, requestedTicker ?? "None"); var macroTask = FetchMacroDataAsync(cancellationToken); string? ticker = requestedTicker; if (string.IsNullOrWhiteSpace(ticker) || string.Equals(ticker.Trim(), cleanIsin, StringComparison.OrdinalIgnoreCase)) { ticker = await _yahooScraper.ResolveTickerFromIsinAsync(cleanIsin, cancellationToken); } var querySymbol = !string.IsNullOrEmpty(ticker) ? ticker : cleanIsin; var (vix, gspc, dxy) = await macroTask; // Lade 2y Daten für saubere Indikator-Aufwärmphasen var yahooResult = await _yahooScraper.FetchHistoricalCandlesWithCurrencyAsync(querySymbol, "2y", "1d", cancellationToken); var candles = yahooResult.Candles; var currency = yahooResult.Currency; if (candles.Count == 0 && querySymbol != cleanIsin) { yahooResult = await _yahooScraper.FetchHistoricalCandlesWithCurrencyAsync(cleanIsin, "2y", "1d", cancellationToken); candles = yahooResult.Candles; currency = yahooResult.Currency; } if (candles.Count == 0) { _logger.LogWarning("[{Channel}] No candles retrieved for {Symbol}", "TechnicalAnalysisChannel", querySymbol); return null; } _candleCache[cleanIsin] = (candles.Select(CloneCandle).ToList(), querySymbol, currency, DateTime.UtcNow); await MergeLivePriceAsync(cleanIsin, candles, querySymbol, currency, cancellationToken); var resultDto = BuildDto(cleanIsin, querySymbol, currency, candles, vix, gspc, dxy); await PersistToDbCacheAsync(cleanIsin, querySymbol, resultDto, cancellationToken); return resultDto; } private async Task BuildAnalysisWithLivePriceAsync( string cleanIsin, List cachedCandles, string querySymbol, string currency, CancellationToken cancellationToken) { var candles = cachedCandles.Select(CloneCandle).ToList(); var livePriceTask = FetchLivePriceAsync(cleanIsin, cancellationToken); var macroTask = FetchMacroDataAsync(cancellationToken); await Task.WhenAll(livePriceTask, macroTask); var (livePrice, liveBid, liveAsk, preChange) = await livePriceTask; // Task-Result direkt nutzen var (vix, gspc, dxy) = await macroTask; ApplyLivePriceToCandles(cleanIsin, candles, querySymbol, currency, livePrice, liveBid, liveAsk); return BuildDto(cleanIsin, querySymbol, currency, candles, vix, gspc, dxy); } private async Task MergeLivePriceAsync(string cleanIsin, List candles, string querySymbol, string currency, CancellationToken cancellationToken) { var (livePrice, liveBid, liveAsk, _) = await FetchLivePriceAsync(cleanIsin, cancellationToken); ApplyLivePriceToCandles(cleanIsin, candles, querySymbol, currency, livePrice, liveBid, liveAsk); } private void ApplyLivePriceToCandles( string cleanIsin, List candles, string querySymbol, string candleCurrency, decimal? livePrice, decimal? liveBid, decimal? liveAsk) { if (!livePrice.HasValue || livePrice.Value <= 0m) return; // Währungsschutz: Trade Republic liefert IMMER EUR. // Wenn die Kerzenhistorie USD ist (z.B. AAPL), darf der EUR-Livepreis NICHT direkt injiziert werden! if (candleCurrency.Equals("USD", StringComparison.OrdinalIgnoreCase) && !cleanIsin.StartsWith("DE") && !cleanIsin.StartsWith("AT")) { _logger.LogDebug("[{Channel}] Skipping direct EUR live price injection for USD asset {Isin}", "TechnicalAnalysisChannel", cleanIsin); return; } var today = DateTime.UtcNow.Date; var lastCandle = candles.LastOrDefault(c => c.Timestamp.Date == today) ?? candles.LastOrDefault(); if (lastCandle != null) { lastCandle.Close = livePrice.Value; lastCandle.High = Math.Max(lastCandle.High, livePrice.Value); lastCandle.Low = Math.Min(lastCandle.Low, livePrice.Value); if (liveBid.HasValue) lastCandle.Bid = liveBid.Value; if (liveAsk.HasValue) lastCandle.Ask = liveAsk.Value; } } private async Task<(decimal? livePrice, decimal? liveBid, decimal? liveAsk, decimal? preChange)> FetchLivePriceAsync( string cleanIsin, CancellationToken cancellationToken) { decimal? livePrice = null; decimal? liveBid = null; decimal? liveAsk = null; decimal? preChange = null; try { using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); cts.CancelAfter(1500); var trTask = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); int? subId = await _trService.SubscribeRealtimeTickerAsync(cleanIsin, tick => { decimal? effectivePrice = tick.Bid?.PriceValue > 0m ? tick.Bid.PriceValue : (tick.Last?.PriceValue > 0m ? tick.Last.PriceValue : null); if (effectivePrice.HasValue) { livePrice = tick.Last?.PriceValue ?? effectivePrice.Value; liveBid = tick.Bid?.PriceValue; liveAsk = tick.Ask?.PriceValue; decimal prePrice = tick.Pre?.PriceValue ?? 0m; if (prePrice > 0m) { preChange = Math.Round(((effectivePrice.Value - prePrice) / prePrice) * 100m, 2); } trTask.TrySetResult(true); } }, cts.Token); if (subId.HasValue) { try { await trTask.Task.WaitAsync(cts.Token); } catch (OperationCanceledException) { } await _trService.UnsubscribeRealtimeTickerAsync(subId.Value); } } catch (Exception ex) { _logger.LogWarning(ex, "[{Channel}] Real-time price fetch skipped for ISIN {Isin}", "TechnicalAnalysisChannel", cleanIsin); } return (livePrice, liveBid, liveAsk, preChange); } private async Task<(MacroDataEntity vix, MacroDataEntity gspc, MacroDataEntity dxy)> FetchMacroDataAsync( CancellationToken cancellationToken) { var vixTask = _yahooScraper.FetchMacroTickerAsync("^VIX", cancellationToken); var gspcTask = _yahooScraper.FetchMacroTickerAsync("^GSPC", cancellationToken); var dxyTask = _yahooScraper.FetchMacroTickerAsync("DX-Y.NY", cancellationToken); await Task.WhenAll(vixTask, gspcTask, dxyTask); var vix = await vixTask ?? new MacroDataEntity { Symbol = "^VIX", Value = 18.5m, TrendState = "Moderate" }; var gspc = await gspcTask ?? new MacroDataEntity { Symbol = "^GSPC", Value = 5500m, TrendState = "Bullish" }; var dxy = await dxyTask ?? new MacroDataEntity { Symbol = "DX-Y.NY", Value = 104.2m, TrendState = "Neutral" }; return (vix, gspc, dxy); } private TechnicalAnalysisDto BuildDto(string cleanIsin, string querySymbol, string currency, List candles, MacroDataEntity vix, MacroDataEntity gspc, MacroDataEntity dxy) { var vixRegime = vix.Value > 25m ? "HighVolatility" : (vix.Value > 18m ? "Moderate" : "LowVolatility"); var summaryText = $"Markt-Vola (VIX: {vix.Value:F1}) ist {vixRegime}. S&P 500 Trend ist {gspc.TrendState}. DXY: {dxy.Value:F1}."; var marketRegime = new MarketRegimeDto( VixValue: vix.Value, VixRegime: vixRegime, MarketTrend: gspc.TrendState, DxyValue: dxy.Value, DxyState: dxy.TrendState == "Bullish" ? "DollarStrengthening" : "DollarWeakening", SummaryText: summaryText); var (indicators, patterns, signals) = _calculator.CalculateAnalysis(candles, currency); var candleDtos = candles.Select(c => new CandleDto( Timestamp: c.Timestamp, Open: c.Open, High: c.High, Low: c.Low, Close: c.Close, Volume: c.Volume, Bid: c.Bid, Ask: c.Ask)).ToList(); return new TechnicalAnalysisDto( Isin: cleanIsin, Ticker: querySymbol, CompanyName: querySymbol, LastUpdated: DateTime.UtcNow, Candles: candleDtos, Indicators: indicators, Patterns: patterns, Signals: signals, MarketRegime: marketRegime, Currency: currency); } private async Task GetFromDbCacheAsync(string cleanIsin, string? requestedTicker, CancellationToken cancellationToken) { try { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var cached = await db.CachedAnalyses .AsNoTracking() .FirstOrDefaultAsync(c => c.Isin == cleanIsin, cancellationToken); if (cached != null && DateTime.UtcNow - cached.CalculatedAt < DbCacheTtl) { if (!string.IsNullOrWhiteSpace(requestedTicker) && !string.Equals(requestedTicker.Trim(), cleanIsin, StringComparison.OrdinalIgnoreCase) && !string.Equals(cached.Ticker, requestedTicker, StringComparison.OrdinalIgnoreCase)) { return null; // Ticker mismatch, force refresh required } return JsonSerializer.Deserialize(cached.AnalysisJson); } } catch (Exception ex) { _logger.LogWarning(ex, "[{Channel}] Failed to read DB cache for ISIN {Isin}", "TechnicalAnalysisChannel", cleanIsin); } return null; } private async Task PersistToDbCacheAsync(string cleanIsin, string querySymbol, TechnicalAnalysisDto dto, CancellationToken cancellationToken) { try { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var json = JsonSerializer.Serialize(dto); var existing = await db.CachedAnalyses.FirstOrDefaultAsync(c => c.Isin == cleanIsin, cancellationToken); if (existing != null) { existing.Ticker = querySymbol; existing.AnalysisJson = json; existing.CalculatedAt = DateTime.UtcNow; } else { db.CachedAnalyses.Add(new CachedAnalysisEntity { Isin = cleanIsin, Ticker = querySymbol, AnalysisJson = json, CalculatedAt = DateTime.UtcNow }); } await db.SaveChangesAsync(cancellationToken); } catch (Exception ex) { _logger.LogError(ex, "[{Channel}] Failed to persist TA DB cache for ISIN {Isin}", "TechnicalAnalysisChannel", cleanIsin); } } private static MarketCandleEntity CloneCandle(MarketCandleEntity c) => new() { Symbol = c.Symbol, Interval = c.Interval, Timestamp = c.Timestamp, Open = c.Open, High = c.High, Low = c.Low, Close = c.Close, Volume = c.Volume, Bid = c.Bid, Ask = c.Ask }; }