using System.IO; using System.Text.Json; using FinlyticCore.Dtos; using FinlyticCore.Dtos.News; using FinlyticCore.Dtos.Sentiment; using FinlyticCore.Models; using FinlyticCore.Util; using FinlyticNews.Entities; using FinlyticNews.Services; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace FinlyticNews.Util; /// /// A managed MQTT client for broadcasting completed news articles and responding to RPC requests. /// public class NewsMqttClient : ManagedMqttClient, IHostedService { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private readonly IConfiguration _configuration; public NewsMqttClient( ILogger logger, IServiceScopeFactory scopeFactory, IConfiguration configuration) : base(logger) { _scopeFactory = scopeFactory; _logger = logger; _configuration = configuration; } /// public async Task StartAsync(CancellationToken cancellationToken) { var config = new MqttConfiguration { Host = _configuration["MQTT:Host"] ?? _configuration["MQTT__Host"] ?? "localhost", Port = Convert.ToInt32(_configuration["MQTT:Port"] ?? _configuration["MQTT__Port"] ?? "1883"), ClientId = $"{(_configuration["MQTT:ClientId"] ?? _configuration["MQTT__ClientId"] ?? "FinlyticNews")}_{Guid.NewGuid()}" }; _logger.LogInformation("Starting News MQTT client. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } /// public async Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("Stopping News MQTT client and disconnecting."); await DisconnectAsync(); } /// protected override async Task OnConnectedAsync() { _logger.LogInformation("News MQTT client connected. Subscribing to RPC topics..."); // ZUSAMMENGELEGT: Unified News Fetching (news_Get deckt news_GetDaily mit ab) await SubscribeAsync("services/request/news_Get/#"); await SubscribeAsync("services/request/news_GetById/#"); await SubscribeAsync("services/request/news_GetPending/#"); await SubscribeAsync("services/request/news_UpdateStatus/#"); await SubscribeAsync("services/request/health_Ping/#"); await SubscribeAsync("services/config/updated/#"); } /// /// Broadcasts a newly processed news article to downstream subscribers. /// public async Task BroadcastArticleAsync(NewsArticleDto article) { // 1. Primärer System-Broadcast const string topic = "services/news/completed"; _logger.LogInformation("Broadcasting completed article to MQTT topic: {Topic}. ID: {Id}", topic, article.Id); await PublishAsync(topic, article); // 2. Zielgerichteter ISIN-Stream für Echtzeit-Frontend-Feeds var firstIsin = article.MatchedAssets.FirstOrDefault()?.Isin; if (!string.IsNullOrWhiteSpace(firstIsin)) { string isinTopic = $"finlytic/news/stream/{firstIsin.Trim().ToLowerInvariant()}"; await PublishAsync(isinTopic, article); } } /// protected override async Task OnMessageReceivedAsync(string topic, string payload) { if (string.IsNullOrWhiteSpace(topic)) return; // 1. System Config Updates if (topic.StartsWith("services/config/updated", StringComparison.OrdinalIgnoreCase)) { if (topic.EndsWith("FinlyticNews", StringComparison.OrdinalIgnoreCase)) { await OnConfigUpdatedAsync(payload); } return; } var segments = topic.Split('/'); if (segments.Length < 4) return; var channel = segments[2]; var correlationId = segments[^1]; // 2. Unified RPC Dispatching via Switch switch (channel) { case "news_Get": case "news_GetDaily": // Abwärtskompatibel weitergeleitet await OnGetNewsAsync(payload, correlationId); break; case "news_GetById": await OnGetNewsByIdAsync(payload, correlationId); break; case "news_GetPending": await OnGetPendingNewsAsync(payload, correlationId); break; case "news_UpdateStatus": await OnUpdateNewsStatusAsync(payload, correlationId); break; case "health_Ping": await OnHealthPingAsync(segments, correlationId); break; default: _logger.LogDebug("Received unhandled RPC channel: {Channel}", channel); break; } } private async Task OnGetNewsAsync(string payload, string correlationId) { _logger.LogInformation("Received RPC news_Get request. Correlation: {CorrelationId}", correlationId); int limit = 20; int offset = 0; string? isin = null; DateTime? date = null; string? status = null; string? searchQuery = null; if (!string.IsNullOrWhiteSpace(payload)) { try { using var doc = JsonDocument.Parse(payload); var root = doc.RootElement; if (root.TryGetProperty("limit", out var limitProp) && limitProp.TryGetInt32(out var parsedLimit)) limit = parsedLimit; if (root.TryGetProperty("offset", out var offsetProp) && offsetProp.TryGetInt32(out var parsedOffset)) offset = parsedOffset; if (root.TryGetProperty("isin", out var isinProp) && isinProp.ValueKind == JsonValueKind.String) isin = isinProp.GetString(); if (root.TryGetProperty("symbol", out var symProp) && symProp.ValueKind == JsonValueKind.String && string.IsNullOrEmpty(isin)) isin = symProp.GetString(); if (root.TryGetProperty("status", out var stProp) && stProp.ValueKind == JsonValueKind.String) status = stProp.GetString(); if (root.TryGetProperty("query", out var qProp) && qProp.ValueKind == JsonValueKind.String) searchQuery = qProp.GetString(); if (root.TryGetProperty("date", out var dProp) && dProp.ValueKind == JsonValueKind.String) { var dStr = dProp.GetString(); if (!string.IsNullOrWhiteSpace(dStr)) { if (string.Equals(dStr, "today", StringComparison.OrdinalIgnoreCase)) date = DateTime.UtcNow.Date; else if (DateTime.TryParse(dStr, out var parsedDate)) date = parsedDate.Date; } } if (root.TryGetProperty("hasSentiment", out var hsProp)) { bool isTrue = hsProp.ValueKind == JsonValueKind.True || (hsProp.ValueKind == JsonValueKind.String && bool.TryParse(hsProp.GetString(), out var b) && b); if (isTrue && string.IsNullOrEmpty(status)) status = "Analyzed"; } } catch (Exception ex) { _logger.LogWarning(ex, "Failed to parse RPC payload on news_Get"); } } try { using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); var articles = await dbService.GetFilteredNewsAsync(limit, offset, isin, date, status, searchQuery); var dtos = (await Task.WhenAll(articles.Select(a => MapToDtoAsync(a)))).ToList(); string responseTopic = $"services/response/news_Get/{correlationId}"; _logger.LogInformation("Publishing RPC response to {ResponseTopic} with {Count} articles.", responseTopic, dtos.Count); await PublishAsync(responseTopic, dtos); } catch (Exception ex) { _logger.LogError(ex, "Failed to compile RPC response for news_Get"); } } private async Task OnGetNewsByIdAsync(string payload, string correlationId) { _logger.LogInformation("Received RPC news_GetById request. Correlation: {CorrelationId}", correlationId); string responseTopic = $"services/response/news_GetById/{correlationId}"; if (string.IsNullOrWhiteSpace(payload)) { await PublishAsync(responseTopic, (object?)null); return; } try { var request = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.ArticleRequest); var targetIdStr = request?.ArticleId ?? request?.Id; if (Guid.TryParse(targetIdStr, out var articleId)) { using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); var article = await dbService.GetArticleByIdAsync(articleId); if (article != null) { var dto = await MapToDtoAsync(article); await PublishAsync(responseTopic, dto); return; } } } catch (Exception ex) { _logger.LogError(ex, "Failed to execute RPC news_GetById"); } await PublishAsync(responseTopic, (object?)null); } private async Task OnGetPendingNewsAsync(string payload, string correlationId) { _logger.LogInformation("Received RPC news_GetPending request. Correlation: {CorrelationId}", correlationId); int limit = 10; if (!string.IsNullOrWhiteSpace(payload)) { try { using var doc = JsonDocument.Parse(payload); if (doc.RootElement.TryGetProperty("limit", out var limitProp) && limitProp.TryGetInt32(out var parsedLimit)) { limit = Math.Min(parsedLimit, 10); } } catch { } } try { using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); var pendingArticles = await dbService.GetArticlesByStatusAsync("Pending"); var dtos = (await Task.WhenAll(pendingArticles.Take(limit).Select(a => MapToDtoAsync(a)))).ToList(); string responseTopic = $"services/response/news_GetPending/{correlationId}"; _logger.LogInformation("Publishing RPC news_GetPending response to {ResponseTopic} with {Count} articles.", responseTopic, dtos.Count); await PublishAsync(responseTopic, dtos); } catch (Exception ex) { _logger.LogError(ex, "Failed to publish RPC pending news response"); } } private async Task OnUpdateNewsStatusAsync(string payload, string correlationId) { _logger.LogInformation("Received RPC news_UpdateStatus request. Correlation: {CorrelationId}", correlationId); UpdateNewsStatusResponse response; try { var request = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.UpdateNewsStatusRequest); if (request != null && request.Id != Guid.Empty) { using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); await dbService.UpdateArticleStatusAsync(request.Id, request.Status); response = new UpdateNewsStatusResponse(true, "Status updated successfully."); } else { response = new UpdateNewsStatusResponse(false, "Invalid payload."); } } catch (Exception ex) { _logger.LogError(ex, "Failed to execute RPC news_UpdateStatus"); response = new UpdateNewsStatusResponse(false, ex.Message); } string responseTopic = $"services/response/news_UpdateStatus/{correlationId}"; await PublishAsync(responseTopic, response); } private async Task OnConfigUpdatedAsync(string payload) { _logger.LogInformation("Received config update event for FinlyticNews."); try { using var doc = JsonDocument.Parse(payload); if (doc.RootElement.TryGetProperty("settings", out var settingsProp)) { var dict = JsonSerializer.Deserialize>(settingsProp.GetRawText()); if (dict != null && dict.Count > 0) { using var scope = _scopeFactory.CreateScope(); var settingsDb = scope.ServiceProvider.GetRequiredService(); await settingsDb.UpdateSettingsFromDictionaryAsync(dict); _logger.LogInformation("Successfully persisted {Count} updated settings.", dict.Count); } } } catch (Exception ex) { _logger.LogError(ex, "Error processing MQTT config update event"); } } private async Task OnHealthPingAsync(string[] segments, string correlationId) { // Zerlegt den Topic-Pfad z. B. services/request/health_Ping/FinlyticNews/{correlationId} bool isForMe = segments.Length >= 5 && segments[3].Equals("FinlyticNews", StringComparison.OrdinalIgnoreCase); if (isForMe) { string respTopic = $"services/response/health_Ping/{correlationId}"; await PublishAsync(respTopic, new ServiceHealthResponse("FinlyticNews", "Online", DateTime.UtcNow, "Connected")); _logger.LogInformation("Responded to live health_Ping RPC request [CorrelationId: {CorrelationId}].", correlationId); } } /// /// Maps a NewsArticleEntity to a NewsArticleDto while performing zero-latency disk lookups for sentiment summaries. /// private async Task MapToDtoAsync(NewsArticleEntity a) { string? sentimentLabel = null; double? sentimentScore = null; double? confidence = null; FinBertResultDto? finbertResult = null; try { var targetId = a.Id.ToString(); IsinAnalysisEntry? sentimentEntry = null; // 1. Snappy Local Disk Check for Article File var articlePath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "articles", $"{targetId}.json"); if (File.Exists(articlePath)) { try { var json = await File.ReadAllTextAsync(articlePath); sentimentEntry = JsonSerializer.Deserialize(json, FinlyticJsonSerializerContext.Default.IsinAnalysisEntry); } catch { } } // 2. ISIN Summary File Fallback if (sentimentEntry == null && a.MatchedAssets != null && a.MatchedAssets.Count > 0) { foreach (var asset in a.MatchedAssets) { if (string.IsNullOrWhiteSpace(asset.Isin)) continue; var isinPath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "isin", $"{asset.Isin.Trim()}.json"); if (File.Exists(isinPath)) { try { var json = await File.ReadAllTextAsync(isinPath); var isinDoc = JsonSerializer.Deserialize(json, FinlyticJsonSerializerContext.Default.IsinSentimentSummaryDto); var match = isinDoc?.Analyses?.FirstOrDefault(entry => string.Equals(entry.Article?.ArticleId?.Trim(), targetId, StringComparison.OrdinalIgnoreCase)); if (match != null) { sentimentEntry = match; break; } } catch { } } } } if (sentimentEntry?.FinbertResult != null) { finbertResult = sentimentEntry.FinbertResult; sentimentLabel = finbertResult.Label; sentimentScore = finbertResult.CompoundScore; confidence = finbertResult.Confidence; } } catch (Exception ex) { _logger.LogTrace(ex, "[MapToDtoAsync] Sentiment fetch skipped for article {Id}", a.Id); } return new NewsArticleDto { Id = a.Id, Title = a.Title, Author = a.Author, Summary = a.Summary, ContentRaw = a.ContentRaw, Language = a.Language, SourceUrl = a.SourceUrl, ScrapedAt = a.ScrapedAt, PublishedAt = a.PublishedAt, Status = a.Status, Sentiment = sentimentLabel, SentimentScore = sentimentScore, Confidence = confidence, FinbertResult = finbertResult, MatchedAssets = (a.MatchedAssets ?? []).Select(m => new MatchedAssetDto { Name = m.Name, Isin = m.Isin }).ToList() }; } }