using System; using System.Collections.Concurrent; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using FinlyticBackend.Database; using FinlyticBackend.Hubs; using FinlyticBackend.Services; using FinlyticCore.Dtos.Logging; using FinlyticCore.Dtos.News; using FinlyticCore.Models; using FinlyticCore.Models.Auth; using FinlyticCore.Models.Trades; using FinlyticCore.Util; using Microsoft.AspNetCore.SignalR; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace FinlyticBackend.Util; /// /// Subscribes to general broadcast MQTT topics and forwards them to SignalR clients and FCM push services. /// public class BackendMqttBridge : ManagedMqttClient, IHostedService { public static readonly ConcurrentDictionary FundamentalsCache = new(StringComparer.OrdinalIgnoreCase); public static readonly ConcurrentDictionary TechnicalsCache = new(StringComparer.OrdinalIgnoreCase); public static readonly ConcurrentDictionary> ServiceLogsRingBuffer = new(StringComparer.OrdinalIgnoreCase); private readonly IConfiguration _configuration; private readonly IServiceScopeFactory _scopeFactory; private readonly IHubContext _hubContext; private readonly IHubContext _tradeHubContext; private readonly IHubContext _newsHubContext; private readonly IHubContext _logHubContext; private readonly IFirebaseNotificationService _firebaseService; private readonly ILogger _logger; public BackendMqttBridge( IConfiguration configuration, IServiceScopeFactory scopeFactory, IHubContext hubContext, IHubContext tradeHubContext, IHubContext newsHubContext, IHubContext logHubContext, IFirebaseNotificationService firebaseService, ILogger logger) : base(logger) { _configuration = configuration; _scopeFactory = scopeFactory; _hubContext = hubContext; _tradeHubContext = tradeHubContext; _newsHubContext = newsHubContext; _logHubContext = logHubContext; _firebaseService = firebaseService; _logger = logger; } /// 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"] ?? "finlytic_backend_bridge")}_{Guid.NewGuid()}" }; _logger.LogInformation("Starting Backend MQTT Bridge. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } /// public async Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("Stopping Backend MQTT Bridge."); await DisconnectAsync(); } /// protected override async Task OnConnectedAsync() { _logger.LogInformation("Backend MQTT Bridge connected. Subscribing to broadcast topics..."); await SubscribeAsync("finlytic/trades/proposed/#"); await SubscribeAsync("finlytic/trades/updates/#"); // News topics await SubscribeAsync("services/news/completed"); await SubscribeAsync("finlytic/news/#"); await SubscribeAsync("finlytic/sentiment/#"); // Fundamentals & Technicals await SubscribeAsync("finlytic/fundamentals/#"); await SubscribeAsync("finlytic/assets/fundamentals/#"); await SubscribeAsync("finlytic/technicalanalysis/#"); await SubscribeAsync("finlytic/ta/#"); // Real-time Logs await SubscribeAsync("finlytic/logs/#"); } /// protected override async Task OnMessageReceivedAsync(string topic, string payloadStr) { if (string.IsNullOrWhiteSpace(topic) || string.IsNullOrWhiteSpace(payloadStr)) return; try { if (topic.StartsWith("finlytic/logs/", StringComparison.OrdinalIgnoreCase)) { await HandleLogMessageAsync(payloadStr); } else if (topic.StartsWith("finlytic/trades/proposed/", StringComparison.OrdinalIgnoreCase) || topic.StartsWith("finlytic/trades/update", StringComparison.OrdinalIgnoreCase)) { await HandleTradeProposalAsync(payloadStr); } else if (topic.StartsWith("finlytic/trades/updates/", StringComparison.OrdinalIgnoreCase)) { await HandleTradeUpdateAsync(payloadStr); } else if (topic.Equals("services/news/completed", StringComparison.OrdinalIgnoreCase) || topic.StartsWith("finlytic/news/", StringComparison.OrdinalIgnoreCase) || topic.StartsWith("finlytic/sentiment/", StringComparison.OrdinalIgnoreCase)) { await HandleNewsArticleAsync(payloadStr); } else if (topic.StartsWith("finlytic/fundamentals/", StringComparison.OrdinalIgnoreCase) || topic.StartsWith("finlytic/assets/fundamentals/", StringComparison.OrdinalIgnoreCase)) { HandleFundamentalsCache(topic, payloadStr); } else if (topic.StartsWith("finlytic/technicalanalysis/", StringComparison.OrdinalIgnoreCase) || topic.StartsWith("finlytic/ta/", StringComparison.OrdinalIgnoreCase)) { HandleTechnicalsCache(topic, payloadStr); } } catch (Exception ex) { _logger.LogError(ex, "Error processing message in Backend MQTT Bridge on topic {Topic}", topic); } } private async Task HandleLogMessageAsync(string payloadStr) { var logDto = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.LogMessageDto); if (logDto == null) return; string serviceKey = logDto.ServiceName; var queue = ServiceLogsRingBuffer.GetOrAdd(serviceKey, _ => new ConcurrentQueue()); queue.Enqueue(logDto); // Keep buffer capped at 250 entries while (queue.Count > 250 && queue.TryDequeue(out _)) { } // Broadcast to SignalR clients await _logHubContext.Clients.Group(serviceKey).SendAsync("ReceiveLogMessage", logDto); await _logHubContext.Clients.All.SendAsync("ReceiveLogMessage", logDto); } private async Task HandleTradeProposalAsync(string payloadStr) { var proposal = JsonSerializer.Deserialize(payloadStr); if (proposal == null) return; await _hubContext.Clients.All.OnTradeProposed(proposal); await _tradeHubContext.Clients.All.SendAsync("ReceiveTradeUpdate", proposal); _logger.LogInformation("Broadcasted Trade Proposal {AnalysisId} via SignalR.", proposal.AnalysisId); using var scope = _scopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); var fcmTokens = await dbContext.UserDeviceTokens.AsNoTracking().Select(t => t.FcmToken).ToListAsync(); if (fcmTokens.Count > 0) { await _firebaseService.SendTradeProposalNotificationAsync(proposal, fcmTokens); } } private async Task HandleTradeUpdateAsync(string payloadStr) { var update = JsonSerializer.Deserialize(payloadStr); if (update == null) return; await _hubContext.Clients.All.OnTradeUpdated(update); await _tradeHubContext.Clients.All.SendAsync("ReceiveTradeUpdate", update); if (string.Equals(update.Recommendation, "Close", StringComparison.OrdinalIgnoreCase) || string.Equals(update.Recommendation, "AdjustSL", StringComparison.OrdinalIgnoreCase)) { using var scope = _scopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); var fcmTokens = await dbContext.UserDeviceTokens.AsNoTracking().Select(t => t.FcmToken).ToListAsync(); if (fcmTokens.Count > 0) { await _firebaseService.SendTradeUpdateNotificationAsync(update, fcmTokens); } } } private async Task HandleNewsArticleAsync(string payloadStr) { var article = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.NewsArticleDto); if (article == null) return; await _newsHubContext.Clients.All.SendAsync("ReceiveNewArticle", article); _logger.LogInformation("Broadcasted live news item '{Title}' over SignalR NewsHub.", article.Title); } private void HandleFundamentalsCache(string topic, string payloadStr) { using var jsonDoc = JsonDocument.Parse(payloadStr); var root = jsonDoc.RootElement.Clone(); string? isin = root.TryGetProperty("isin", out var isinProp) ? isinProp.GetString() : null; string? ticker = root.TryGetProperty("primaryTicker", out var tickerProp) ? tickerProp.GetString() : (root.TryGetProperty("ticker", out var tProp) ? tProp.GetString() : null); if (!string.IsNullOrEmpty(isin)) FundamentalsCache[isin] = root; if (!string.IsNullOrEmpty(ticker)) FundamentalsCache[ticker] = root; string topicKey = topic.Split('/').LastOrDefault() ?? string.Empty; if (!string.IsNullOrEmpty(topicKey)) FundamentalsCache[topicKey] = root; _logger.LogInformation("Cached fundamentals payload from MQTT topic {Topic}.", topic); } private void HandleTechnicalsCache(string topic, string payloadStr) { using var jsonDoc = JsonDocument.Parse(payloadStr); var root = jsonDoc.RootElement.Clone(); string? symbol = root.TryGetProperty("symbol", out var sProp) ? sProp.GetString() : null; if (!string.IsNullOrEmpty(symbol)) TechnicalsCache[symbol] = root; string topicKey = topic.Split('/').LastOrDefault() ?? string.Empty; if (!string.IsNullOrEmpty(topicKey)) TechnicalsCache[topicKey] = root; _logger.LogInformation("Cached technical analysis payload from MQTT topic {Topic}.", topic); } }