using System; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using FinlyticAssets.Entities; using FinlyticAssets.Services; using FinlyticCore.Dtos.Settings; using FinlyticCore.Models; using FinlyticCore.Services; using FinlyticCore.Util; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace FinlyticAssets.Util; /// /// Represents a managed MQTT client acting as a server-side RPC provider within the asset microservice. /// public class AssetsMqttClient : ManagedMqttClient, IHostedService { private readonly ILogger _logger; private readonly IServiceScopeFactory _scopeFactory; private readonly IConfiguration _configuration; public AssetsMqttClient( ILogger logger, IServiceScopeFactory scopeFactory, IConfiguration configuration) : base(logger) { _logger = logger; _scopeFactory = scopeFactory; _configuration = configuration; } /// /// Starts the MQTT client and connects to the configured broker. /// 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"] ?? "FinlyticAssets")}_{Guid.NewGuid()}" }; _logger.LogInformation("Starting Assets MQTT client. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } /// /// Gracefully stops and disconnects the MQTT client. /// public async Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("Stopping Assets MQTT client."); await DisconnectAsync(); } /// /// Invoked automatically once the connection to the MQTT broker is successfully established or restored. /// protected override async Task OnConnectedAsync() { _logger.LogInformation("Assets MQTT Client connected. Subscribing to topics..."); await SubscribeAsync("services/request/assets_Get/#"); await SubscribeAsync("services/request/assets_Search/#"); await SubscribeAsync("services/request/assets_GetDiscovery/#"); await SubscribeAsync("services/request/assets_GetDerivatives/#"); await SubscribeAsync("services/request/assets_FetchLogo/#"); await SubscribeAsync("services/request/assets_settings_GetAll/#"); await SubscribeAsync("services/request/assets_settings_Update/#"); await SubscribeAsync("services/request/health_Ping/#"); await SubscribeAsync("services/config/updated/#"); FinlyticCore.Services.FinlyticLogBroadcaster.OnLogPublished = async (logDto) => { if (IsConnected && string.Equals(logDto.ServiceName, "FinlyticAssets", StringComparison.OrdinalIgnoreCase)) { await PublishAsync("finlytic/logs/FinlyticAssets", logDto); } }; } /// /// Processes incoming messages on the subscribed topics. /// protected override async Task OnMessageReceivedAsync(string topic, string payload) { if (topic.StartsWith("services/config/updated", StringComparison.OrdinalIgnoreCase)) { await HandleConfigUpdatedAsync(topic, payload); return; } var segments = topic.Split('/'); if (segments.Length < 4) return; var channel = segments[2]; var correlationId = segments[segments.Length - 1]; if (topic.Contains("health_Ping", StringComparison.OrdinalIgnoreCase)) { await HandleHealthPingAsync(topic, segments, correlationId); return; } if (topic.StartsWith("services/request/assets_settings_GetAll", StringComparison.OrdinalIgnoreCase)) { await HandleSettingsGetAllAsync(correlationId); return; } if (topic.StartsWith("services/request/assets_settings_Update", StringComparison.OrdinalIgnoreCase)) { await HandleSettingsUpdateAsync(payload, correlationId); return; } try { using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); var indexService = scope.ServiceProvider.GetRequiredService(); if (channel == "assets_FetchLogo") { await HandleFetchLogoAsync(payload, correlationId, indexService); return; } List responseData = []; switch (channel) { case "assets_Get": responseData = await HandleAssetsGetAsync(payload, dbService); break; case "assets_Search": responseData = await HandleAssetsSearchAsync(payload, dbService); break; case "assets_GetDiscovery": responseData = await HandleAssetsGetDiscoveryAsync(payload, dbService); break; case "assets_GetDerivatives": responseData = (await HandleAssetsGetDerivativesAsync(payload, dbService)).Cast().ToList(); break; } string defaultResponseTopic = $"services/response/{channel}/{correlationId}"; await PublishAsync(defaultResponseTopic, responseData.ToDtoList()); } catch (Exception ex) { OnError(ex); } } private async Task HandleSettingsGetAllAsync(string correlationId) { using var scope = _scopeFactory.CreateScope(); var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var settingsService = scope.ServiceProvider.GetRequiredService(); await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticAssets] [Settings_GetAll] Retrieving all dynamic settings via reflection [CorrelationId: {CorrelationId}]", correlationId); try { var settings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); var responseTopic = $"services/response/assets_settings_GetAll/{correlationId}"; await PublishAsync(responseTopic, settings); await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticAssets] [Settings_GetAll] Published {Count} settings to '{ResponseTopic}'", settings.Count, responseTopic); } catch (Exception ex) { await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticAssets] [Settings_GetAll] Failed to retrieve settings."); } } private async Task HandleSettingsUpdateAsync(string payload, string correlationId) { if (string.IsNullOrWhiteSpace(payload)) return; using var scope = _scopeFactory.CreateScope(); var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var settingsService = scope.ServiceProvider.GetRequiredService(); await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticAssets] [Settings_Update] Processing settings update RPC [CorrelationId: {CorrelationId}]", correlationId); try { Dictionary? updates = null; try { updates = JsonSerializer.Deserialize>(payload); } catch { var list = JsonSerializer.Deserialize>(payload); if (list != null) { updates = new Dictionary(); foreach (var item in list) updates[item.Key] = item.Value; } } if (updates != null && updates.Count > 0) { await settingsService.UpdateSettingsAsync(updates); await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticAssets] [Settings_Update] Successfully updated {Count} settings in database and cache.", updates.Count); } var currentSettings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); var responseTopic = $"services/response/assets_settings_Update/{correlationId}"; await PublishAsync(responseTopic, currentSettings); } catch (Exception ex) { await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticAssets] [Settings_Update] Failed to update settings."); } } private async Task HandleConfigUpdatedAsync(string topic, string payload) { if (!topic.EndsWith("FinlyticAssets", StringComparison.OrdinalIgnoreCase)) return; try { var updatePayload = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.ServiceConfigUpdatePayload); if (updatePayload?.Settings != null && updatePayload.Settings.Count > 0) { using var scope = _scopeFactory.CreateScope(); var settings = scope.ServiceProvider.GetRequiredService(); var dict = updatePayload.Settings.ToDictionary(k => k.Key, v => (object?)v.Value); await settings.UpdateSettingsAsync(dict); } } catch (Exception ex) { using var scope = _scopeFactory.CreateScope(); var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); await finlyticLogger.LogErrorAsync(SettingKeys.AssetsChannel, ex, "[AssetsMqttClient] Error processing MQTT config update event."); } } private async Task HandleHealthPingAsync(string topic, string[] segments, string correlationId) { bool isForMe = segments.Length >= 5 ? segments[3].Equals("FinlyticAssets", StringComparison.OrdinalIgnoreCase) : topic.Contains("FinlyticAssets", StringComparison.OrdinalIgnoreCase); if (isForMe) { string respTopic = $"services/response/health_Ping/{correlationId}"; await PublishAsync(respTopic, new FinlyticCore.Dtos.ServiceHealthResponse("FinlyticAssets", "Online", DateTime.UtcNow, "Connected")); using var scope = _scopeFactory.CreateScope(); var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); await finlyticLogger.LogInfoAsync(SettingKeys.HealthPingChannel, "[AssetsMqttClient] Responded to live health_Ping RPC request [CorrelationId: {CorrelationId}].", correlationId); } } private async Task HandleFetchLogoAsync(string payload, string correlationId, IAssetsIndexService indexService) { var req = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.IsinRequest); string? isin = req?.Isin; string? savedPath = null; if (!string.IsNullOrEmpty(isin)) { savedPath = await indexService.DownloadAndSaveLogoAsync(isin); } string responseTopic = $"services/response/assets_FetchLogo/{correlationId}"; await PublishAsync(responseTopic, new FinlyticCore.Dtos.FetchLogoResponse(isin, savedPath, savedPath != null)); } private async Task> HandleAssetsGetAsync(string payload, IAssetsDbService dbService) { var validReq = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.GetValidAssetRequest); if (validReq != null) { return await dbService.GetValidAssetsByIsinAsync(validReq.Isin); } return []; } private async Task> HandleAssetsSearchAsync(string payload, IAssetsDbService dbService) { var searchReq = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.SearchAssetsRequest); if (searchReq != null) { return await dbService.FindAffectedActiveAssetsAsync(searchReq.SearchQuery); } return []; } private async Task> HandleAssetsGetDiscoveryAsync(string payload, IAssetsDbService dbService) { int limit = 15; if (!string.IsNullOrWhiteSpace(payload)) { try { var discReq = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.GetDiscoveryAssetsRequest); if (discReq != null && discReq.Limit > 0) limit = discReq.Limit; } catch { } } return await dbService.GetDiscoveryAssetsAsync(limit); } private async Task> HandleAssetsGetDerivativesAsync(string payload, IAssetsDbService dbService) { if (string.IsNullOrWhiteSpace(payload)) return []; try { var req = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.GetDerivativesRequest); if (req != null && !string.IsNullOrEmpty(req.UnderlyingIsin)) { return await dbService.GetDerivativesByUnderlyingAsync( req.UnderlyingIsin, req.OptionType, req.TargetLeverage, req.After, req.Page, req.ShouldForceRefresh); } } catch (Exception ex) { using var scope = _scopeFactory.CreateScope(); var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); await finlyticLogger.LogErrorAsync(SettingKeys.AssetsChannel, ex, "[AssetsMqttClient] Error parsing GetDerivativesRequest payload."); } return []; } }