using System; using System.Threading; using System.Threading.Tasks; using FinlyticCore.Models; using FinlyticCore.Util; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace FinlyticBackend.Util; /// /// Managed MQTT client for web gateway endpoints enabling RPC communication with background microservices. /// public class WebMqttClient : ManagedMqttClient, IHostedService { private readonly ILogger _logger; private readonly IConfiguration _configuration; public WebMqttClient(ILogger logger, IConfiguration configuration) : base(logger) { _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"] ?? "finlytic_backend_rpc")}_{Guid.NewGuid()}" }; _logger.LogInformation("Starting Web MQTT RPC Client. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } public async Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("Stopping Web MQTT RPC Client."); await DisconnectAsync(); } protected override async Task OnConnectedAsync() { _logger.LogInformation("Web MQTT RPC client connected. Subscribing to RPC response channels..."); await SubscribeAsync("services/response/news_Get/#"); await SubscribeAsync("services/response/news_GetDaily/#"); await SubscribeAsync("services/response/sentiment_GetArticle/#"); await SubscribeAsync("services/response/sentiment_GetIsin/#"); await SubscribeAsync("services/response/fundamentals_Get/#"); await SubscribeAsync("services/response/events_GetAll/#"); await SubscribeAsync("services/response/events_GetByMonth/#"); await SubscribeAsync("services/response/ta_GetAnalysis/#"); await SubscribeAsync("services/response/tr_GetLivePrice/#"); await SubscribeAsync("services/response/assets_Get/#"); await SubscribeAsync("services/response/assets_Search/#"); await SubscribeAsync("services/response/assets_GetDiscovery/#"); await SubscribeAsync("services/response/trades_Get/#"); await SubscribeAsync("services/response/trades_Close/#"); await SubscribeAsync("services/response/analyzer_TriggerManual/#"); await SubscribeAsync("services/response/health_Ping/#"); } protected override Task OnMessageReceivedAsync(string topic, string payload) { return Task.CompletedTask; } }