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 FinlyticBackend.Database; using FinlyticBackend.Hubs; using FinlyticBackend.Services; using FinlyticBackend.Settings; using FinlyticCore.Dtos; using FinlyticCore.Dtos.Bot; using FinlyticCore.Dtos.Logging; using FinlyticCore.Dtos.News; using FinlyticCore.Dtos.Trading; using FinlyticCore.Models; using FinlyticCore.Services; 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; /// /// Central Managed MQTT Bridge & RPC Gateway for FinlyticBackend. /// 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> ServiceLogsRingBuffer = new(StringComparer.OrdinalIgnoreCase); private readonly IConfiguration _configuration; private readonly IServiceScopeFactory _scopeFactory; private readonly IHubContext _tradeStreamHubContext; private readonly IHubContext _newsHubContext; private readonly IHubContext _logHubContext; private readonly IFirebaseNotificationService _firebaseService; private readonly ILogger _logger; private readonly IFinlyticLogger _finlyticLogger; public BackendMqttBridge( IConfiguration configuration, IServiceScopeFactory scopeFactory, IHubContext tradeStreamHubContext, IHubContext newsHubContext, IHubContext logHubContext, IFirebaseNotificationService firebaseService, ILogger logger, IFinlyticLogger finlyticLogger) : base(logger) { _configuration = configuration; _scopeFactory = scopeFactory; _tradeStreamHubContext = tradeStreamHubContext; _newsHubContext = newsHubContext; _logHubContext = logHubContext; _firebaseService = firebaseService; _logger = logger; _finlyticLogger = finlyticLogger; } /// public async Task StartAsync(CancellationToken cancellationToken) { var config = MqttConfiguration.FromConfiguration(_configuration, "finlytic_backend_gateway"); _logger.LogInformation("Starting Backend MQTT Gateway Bridge. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } /// public async Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("Stopping Backend MQTT Gateway Bridge."); await DisconnectAsync(); } /// protected override async Task OnConnectedAsync() { await _finlyticLogger.LogInfoAsync(BackendSettingKeys.MqttChannel, "[BackendMqttBridge] Backend MQTT Gateway Bridge connected. Subscribing to broadcast topics..."); // RPC response stream await SubscribeAsync(MqttTopics.ResponseWildcard); // Engine & Bot streams await SubscribeAsync(MqttTopics.EngineWildcard); await SubscribeAsync(MqttTopics.BotWildcard); // News & Sentiment streams await SubscribeAsync(MqttTopics.NewsCompleted); await SubscribeAsync(MqttTopics.NewsStreamWildcard); await SubscribeAsync(MqttTopics.SentimentWildcard); // Real-time Logs await SubscribeAsync(MqttTopics.LogsWildcard); // Universe RPC Endpoint for FinlyticTechnicals await SubscribeRpcAsync>( MqttTopics.RequestFilter(MqttTopics.Channels.BackendGetAggregatedFavorites), HandleGetAggregatedFavoritesRpcAsync); // Username lookup RPC Endpoint for FinlyticNotify await SubscribeRpcAsync( MqttTopics.RequestFilter(MqttTopics.Channels.BackendGetUsername), HandleGetUsernameRpcAsync); // FinlyticBackend previously never broadcast its OWN structured logs at all (it only relayed other // services' logs received on MqttTopics.LogsWildcard, subscribed above) - it never used // IFinlyticLogger, so FinlyticLogBroadcaster.Broadcast was never invoked for anything happening // inside FinlyticBackend itself. This hook publishes FinlyticBackend's own logs onto the exact same // finlytic/logs/FinlyticBackend topic every other service publishes to, which the wildcard // subscription above then picks straight back up and relays into LogStreamHub like any other // service's messages - no separate "local echo" mechanism needed. FinlyticCore.Services.FinlyticLogBroadcaster.OnLogPublished = async (logDto) => { if (IsConnected && string.Equals(logDto.ServiceName, "FinlyticBackend", StringComparison.OrdinalIgnoreCase)) { await PublishAsync(MqttTopics.Logs("FinlyticBackend"), logDto); } }; } /// protected override async Task OnMessageReceivedAsync(string topic, string payloadStr) { if (string.IsNullOrWhiteSpace(topic) || string.IsNullOrWhiteSpace(payloadStr)) return; try { if (topic.StartsWith(MqttTopics.LogsPrefix, StringComparison.OrdinalIgnoreCase)) { await HandleLogMessageAsync(payloadStr); } else if (topic.StartsWith(MqttTopics.EngineProposalsPrefix, StringComparison.OrdinalIgnoreCase)) { await HandleEngineProposalAsync(payloadStr); } else if (topic.StartsWith(MqttTopics.EngineTradesPrefix, StringComparison.OrdinalIgnoreCase)) { await HandleEngineTradeUpdateAsync(payloadStr); } else if (topic.StartsWith(MqttTopics.BotTradesPrefix, StringComparison.OrdinalIgnoreCase)) { await HandleBotTradeUpdateAsync(payloadStr); } else if (topic.Equals(MqttTopics.NewsCompleted, StringComparison.OrdinalIgnoreCase) || topic.StartsWith(MqttTopics.NewsPrefix, StringComparison.OrdinalIgnoreCase) || topic.StartsWith(MqttTopics.SentimentPrefix, StringComparison.OrdinalIgnoreCase)) { await HandleNewsArticleAsync(payloadStr); } } catch (Exception ex) { await _finlyticLogger.LogErrorAsync(BackendSettingKeys.BackendChannel, ex, "[BackendMqttBridge] Error processing message 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 once to all connected SignalR admin logs clients await _logHubContext.Clients.All.SendAsync("ReceiveLogMessage", logDto); } private async Task HandleEngineProposalAsync(string payloadStr) { var proposal = JsonSerializer.Deserialize(payloadStr); if (proposal == null) return; await _tradeStreamHubContext.Clients.All.ReceiveTradeProposal(proposal); if (!string.IsNullOrWhiteSpace(proposal.UnderlyingIsin)) { await _tradeStreamHubContext.Clients.Group(proposal.UnderlyingIsin.ToUpperInvariant()).ReceiveTradeProposal(proposal); } await _finlyticLogger.LogInfoAsync(BackendSettingKeys.BackendChannel, "[BackendMqttBridge] Broadcasted FinlyticEngine Trade Proposal {ProposalId} for {Isin} via SignalR TradeStreamHub.", proposal.ProposalId, proposal.UnderlyingIsin); // FCM Push Notification for High Score Proposals (Score >= 75) if (proposal.CompositeScore >= 75) { try { 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); } } catch (Exception ex) { await _finlyticLogger.LogWarningAsync(BackendSettingKeys.PushNotificationChannel, ex, "[BackendMqttBridge] Failed to send FCM push notification for proposal {ProposalId}.", proposal.ProposalId); } } } private async Task HandleEngineTradeUpdateAsync(string payloadStr) { var trade = JsonSerializer.Deserialize(payloadStr); if (trade == null) return; await _tradeStreamHubContext.Clients.All.ReceiveTradeUpdate(trade); if (!string.IsNullOrWhiteSpace(trade.UnderlyingIsin)) { await _tradeStreamHubContext.Clients.Group(trade.UnderlyingIsin.ToUpperInvariant()).ReceiveTradeUpdate(trade); } await _finlyticLogger.LogInfoAsync(BackendSettingKeys.BackendChannel, "[BackendMqttBridge] Broadcasted FinlyticEngine Trade Update {TradeId} (Status: {Status}) via SignalR TradeStreamHub.", trade.TradeId, trade.Status); // Push notification on important state transitions (e.g. stopped out or TP2 hit) if (trade.Status == TradeStatus.StoppedOut || trade.Status == TradeStatus.Tp2Hit || trade.Status == TradeStatus.Closed) { try { 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(trade, fcmTokens); } } catch (Exception ex) { await _finlyticLogger.LogWarningAsync(BackendSettingKeys.PushNotificationChannel, ex, "[BackendMqttBridge] Failed to send FCM push notification for trade update {TradeId}.", trade.TradeId); } } } private async Task HandleBotTradeUpdateAsync(string payloadStr) { var botTrade = JsonSerializer.Deserialize(payloadStr, DefaultJsonOptions); if (botTrade != null) { await _tradeStreamHubContext.Clients.All.ReceiveBotPositionUpdate(botTrade); if (!string.IsNullOrWhiteSpace(botTrade.Isin)) { await _tradeStreamHubContext.Clients.Group(botTrade.Isin.ToUpperInvariant()).ReceiveBotPositionUpdate(botTrade); } await _finlyticLogger.LogInfoAsync(BackendSettingKeys.BackendChannel, "[BackendMqttBridge] Broadcasted Bot trade update for {Isin} ({Symbol}, Status: {Status}) via SignalR TradeStreamHub.", botTrade.Isin, botTrade.Symbol, botTrade.Status); } } 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); await _finlyticLogger.LogInfoAsync(BackendSettingKeys.BackendChannel, "[BackendMqttBridge] Broadcasted live news item '{Title}' over SignalR NewsHub.", article.Title); } private async Task> HandleGetAggregatedFavoritesRpcAsync(object? _, string correlationId) { using var scope = _scopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); var isins = await dbContext.UserFavoriteAssets .AsNoTracking() .Select(f => f.Isin) .Where(isin => !string.IsNullOrWhiteSpace(isin)) .Distinct() .ToListAsync(); await _finlyticLogger.LogInfoAsync(BackendSettingKeys.MqttChannel, "[BackendMqttBridge] Responded to backend_GetAggregatedFavorites with {Count} unique ISINs. [CorrelationId: {CorrelationId}]", isins.Count, correlationId); return isins; } private async Task HandleGetUsernameRpcAsync(UserIdRequest? req, string correlationId) { if (req == null || req.UserId == Guid.Empty) return null; using var scope = _scopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); var user = await dbContext.Users .AsNoTracking() .FirstOrDefaultAsync(u => u.Id == req.UserId); if (user == null) return null; string username = !string.IsNullOrWhiteSpace(user.FullName) ? user.FullName.Trim().ToLowerInvariant().Replace(" ", "_") : user.Email.Split('@')[0].Trim().ToLowerInvariant(); return username; } }