Files
Finlytic/FinlyticSimulation/Util/SimulationMqttClient.cs
T

182 lines
9.3 KiB
C#

using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using FinlyticCore.Dtos;
using FinlyticCore.Dtos.Settings;
using FinlyticCore.Dtos.Simulation;
using FinlyticCore.Models;
using FinlyticCore.Services;
using FinlyticCore.Util;
using FinlyticSimulation.Services;
using FinlyticSimulation.Services.Mqtt;
using FinlyticSimulation.Settings;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace FinlyticSimulation.Util;
public class SimulationMqttClient : ManagedMqttClient, IHostedService, ISimulationRpcClient
{
private readonly IConfiguration _configuration;
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger<SimulationMqttClient> _logger;
public SimulationMqttClient(
ILogger<SimulationMqttClient> logger,
IConfiguration configuration,
IServiceScopeFactory scopeFactory) : base(logger)
{
_logger = logger;
_configuration = configuration;
_scopeFactory = scopeFactory;
}
public async Task StartAsync(CancellationToken cancellationToken)
{
var config = MqttConfiguration.FromConfiguration(_configuration, "FinlyticSimulation");
_logger.LogInformation("Starting FinlyticSimulation MQTT client. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId);
await ConnectAsync(config);
}
public async Task StopAsync(CancellationToken cancellationToken)
{
_logger.LogInformation("Stopping FinlyticSimulation MQTT client.");
await DisconnectAsync();
}
protected override async Task OnConnectedAsync()
{
_logger.LogInformation("FinlyticSimulation MQTT client connected. Registering RPC endpoints...");
await SubscribeAsync(MqttTopics.ResponseWildcard);
await SubscribeRpcAsync<BacktestRequestDto, BacktestReportDto>(MqttTopics.RequestFilter(MqttTopics.Channels.SimRunBacktest), HandleRunBacktestRpcAsync);
await SubscribeRpcAsync<GetReliabilityRequest, StrategyAssetReliabilityDto?>(MqttTopics.RequestFilter(MqttTopics.Channels.SimGetReliability), HandleGetReliabilityRpcAsync);
await SubscribeRpcAsync<IsinRequest, List<StrategyAssetReliabilityDto>>(MqttTopics.RequestFilter(MqttTopics.Channels.SimGetMatrixForAsset), HandleGetMatrixRpcAsync);
await SubscribeRpcAsync<GetBacktestHistoryRequest, List<BacktestHistoryEntryDto>>(MqttTopics.RequestFilter(MqttTopics.Channels.SimGetBacktestHistory), HandleGetBacktestHistoryRpcAsync);
await SubscribeRpcAsync<GetBacktestRunDetailRequest, BacktestReportDto?>(MqttTopics.RequestFilter(MqttTopics.Channels.SimGetBacktestRunDetail), HandleGetBacktestRunDetailRpcAsync);
await SubscribeRpcAsync<GetStrategyParametersRequest, StrategyParameterProfileDto?>(MqttTopics.RequestFilter(MqttTopics.Channels.SimGetStrategyParameters), HandleGetStrategyParametersRpcAsync);
await SubscribeRpcAsync<SaveStrategyParametersRequest, StrategyParameterProfileDto>(MqttTopics.RequestFilter(MqttTopics.Channels.SimSaveStrategyParameters), HandleSaveStrategyParametersRpcAsync);
await SubscribeRpcAsync<object, List<DynamicSettingDto>>(MqttTopics.RequestFilter(MqttTopics.Channels.SimSettingsGetAll), HandleSettingsGetAllRpcAsync);
await SubscribeRpcAsync<Dictionary<string, object?>, List<DynamicSettingDto>>(MqttTopics.RequestFilter(MqttTopics.Channels.SimSettingsUpdate), HandleSettingsUpdateRpcAsync);
await SubscribeAsync<object>(MqttTopics.RequestFilter(MqttTopics.Channels.HealthPing), HandleHealthPingRpcAsync);
FinlyticLogBroadcaster.OnLogPublished = async (logDto) =>
{
if (IsConnected && string.Equals(logDto.ServiceName, "FinlyticSimulation", StringComparison.OrdinalIgnoreCase))
{
await PublishAsync(MqttTopics.Logs("FinlyticSimulation"), logDto);
}
};
}
private async Task<BacktestReportDto> HandleRunBacktestRpcAsync(BacktestRequestDto? req, string correlationId)
{
if (req == null) throw new ArgumentNullException(nameof(req));
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
var logger = scope.ServiceProvider.GetRequiredService<IFinlyticLogger<SimulationMqttClient>>();
await logger.LogInfoAsync(SimulationSettingKeys.SimulationChannel,
"[SimulationMqttClient] Processing RPC sim_RunBacktest for {Isin} ({Strategy}) [CorrelationId: {CorrelationId}]",
req.Isin, req.StrategyKey, correlationId);
return await simEngine.RunBacktestAsync(req);
}
private async Task<StrategyAssetReliabilityDto?> HandleGetReliabilityRpcAsync(GetReliabilityRequest? req, string correlationId)
{
if (req == null || string.IsNullOrWhiteSpace(req.Isin)) return null;
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.GetStrategyReliabilityAsync(req.Isin, req.StrategyKey, req.Timeframe);
}
private async Task<List<StrategyAssetReliabilityDto>> HandleGetMatrixRpcAsync(IsinRequest? req, string correlationId)
{
if (req == null || string.IsNullOrWhiteSpace(req.Isin)) return [];
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.GetMatrixForAssetAsync(req.Isin);
}
private async Task<List<BacktestHistoryEntryDto>> HandleGetBacktestHistoryRpcAsync(GetBacktestHistoryRequest? req, string correlationId)
{
if (req == null || string.IsNullOrWhiteSpace(req.Isin)) return [];
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.GetBacktestHistoryAsync(req);
}
private async Task<BacktestReportDto?> HandleGetBacktestRunDetailRpcAsync(GetBacktestRunDetailRequest? req, string correlationId)
{
if (req == null || req.RunId == Guid.Empty) return null;
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.GetBacktestRunDetailAsync(req.RunId);
}
private async Task<StrategyParameterProfileDto?> HandleGetStrategyParametersRpcAsync(GetStrategyParametersRequest? req, string correlationId)
{
if (req == null || string.IsNullOrWhiteSpace(req.Isin) || string.IsNullOrWhiteSpace(req.StrategyKey)) return null;
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.GetStrategyParametersAsync(req.Isin, req.StrategyKey);
}
private async Task<StrategyParameterProfileDto> HandleSaveStrategyParametersRpcAsync(SaveStrategyParametersRequest? req, string correlationId)
{
if (req == null || string.IsNullOrWhiteSpace(req.Isin) || string.IsNullOrWhiteSpace(req.StrategyKey))
{
throw new ArgumentException("Isin and StrategyKey are required to save a parameter profile.");
}
using var scope = _scopeFactory.CreateScope();
var simEngine = scope.ServiceProvider.GetRequiredService<IQuantSimulationEngine>();
return await simEngine.SaveStrategyParametersAsync(req.Isin, req.StrategyKey, req.Parameters ?? new Dictionary<string, decimal>());
}
private async Task<List<DynamicSettingDto>> HandleSettingsGetAllRpcAsync(object? _, string correlationId)
{
using var scope = _scopeFactory.CreateScope();
var settingsService = scope.ServiceProvider.GetRequiredService<ISettingsService>();
return await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SimulationSettingKeys) });
}
private async Task<List<DynamicSettingDto>> HandleSettingsUpdateRpcAsync(Dictionary<string, object?>? updates, string correlationId)
{
using var scope = _scopeFactory.CreateScope();
var settingsService = scope.ServiceProvider.GetRequiredService<ISettingsService>();
if (updates != null && updates.Count > 0)
{
await settingsService.UpdateSettingsAsync(updates);
}
return await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SimulationSettingKeys) });
}
private async Task HandleHealthPingRpcAsync(object? _, string topic, string correlationId)
{
if (topic.Contains("FinlyticSimulation", StringComparison.OrdinalIgnoreCase) || !topic.Contains("/", StringComparison.OrdinalIgnoreCase))
{
string respTopic = MqttTopics.ResponseTopic(MqttTopics.Channels.HealthPing, correlationId);
await PublishAsync(respTopic, new ServiceHealthResponse("FinlyticSimulation", "Online", DateTime.UtcNow, "Connected"));
using var scope = _scopeFactory.CreateScope();
var logger = scope.ServiceProvider.GetRequiredService<IFinlyticLogger<SimulationMqttClient>>();
await logger.LogInfoAsync(SimulationSettingKeys.HealthPingChannel,
"[FinlyticSimulation] Responded to health_Ping RPC [CorrelationId: {CorrelationId}]", correlationId);
}
}
}