From f18f75c1abc534965f5b9ca886c8d413c97cebae Mon Sep 17 00:00:00 2001 From: Kleidukos Date: Sat, 15 Aug 2026 21:30:56 +0200 Subject: [PATCH] feat(backend): generic settings RPC bridge, LogStreamHub SignalR, and log ringbuffer --- .../Controllers/AdminSettingsController.cs | 170 +++++++++++------- FinlyticBackend/Hubs/LogStreamHub.cs | 49 +++++ FinlyticBackend/Program.cs | 4 + FinlyticBackend/Util/BackendMqttBridge.cs | 30 +++- FinlyticBackend/Util/WebMqttClient.cs | 19 +- 5 files changed, 207 insertions(+), 65 deletions(-) create mode 100644 FinlyticBackend/Hubs/LogStreamHub.cs diff --git a/FinlyticBackend/Controllers/AdminSettingsController.cs b/FinlyticBackend/Controllers/AdminSettingsController.cs index 88215bd..499994d 100644 --- a/FinlyticBackend/Controllers/AdminSettingsController.cs +++ b/FinlyticBackend/Controllers/AdminSettingsController.cs @@ -6,6 +6,7 @@ using System.Text.Json.Serialization; using System.Threading.Tasks; using FinlyticBackend.Util; using FinlyticCore.Dtos; +using FinlyticCore.Dtos.Settings; using FinlyticCore.Util; using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.Cors; @@ -58,38 +59,18 @@ public class AdminSettingsController : ControllerBase private readonly WebMqttClient _mqttClient; private readonly ILogger _logger; - // In-memory static store for default UI configuration templates - private static readonly ConcurrentDictionary> _inMemorySettings = new(StringComparer.OrdinalIgnoreCase); - - static AdminSettingsController() + private static readonly Dictionary ServiceRpcPrefixes = new(StringComparer.OrdinalIgnoreCase) { - // Default-Templates für Microservice-Konfigurationen initialisieren - _inMemorySettings["FinlyticAnalyzer"] = new List - { - new("MinSignalScore", "75.0", "double", "Mindest-Score für KI-Trade-Proposals (0-100)", DateTime.UtcNow), - new("VixPanicThreshold", "30.0", "double", "VIX-Wert ab dem Panik-Modus aktiviert wird", DateTime.UtcNow), - new("ProposalTtlMinutes", "180", "int", "Gültigkeitsdauer von Trade-Proposals in Minuten", DateTime.UtcNow) - }; + ["FinlyticFundamentals"] = "fundamentals", + ["FinlyticNews"] = "news", + ["FinlyticTechnicalAnalysis"] = "ta", + ["FinlyticSentiment"] = "sentiment", + ["FinlyticAnalyzer"] = "analyzer", + ["FinlyticTrades"] = "trades", + ["FinlyticAssets"] = "assets" + }; - _inMemorySettings["FinlyticNews"] = new List - { - new("ScrapeIntervalMinutes", "15", "int", "Intervall für das Scraping neuer Nachrichten", DateTime.UtcNow), - new("FinBertBatchSize", "8", "int", "Batch-Größe für die Sentiment-Analyse", DateTime.UtcNow) - }; - - _inMemorySettings["FinlyticTechnicalAnalysis"] = new List - { - new("EmaShortPeriod", "20", "int", "Kurze Periode für EMA-Berechnungen", DateTime.UtcNow), - new("EmaLongPeriod", "50", "int", "Lange Periode für EMA-Berechnungen", DateTime.UtcNow), - new("RsiPeriod", "14", "int", "Standard-Periode für RSI-Berechnung", DateTime.UtcNow) - }; - - _inMemorySettings["FinlyticTrades"] = new List - { - new("ExportFeedbackIntervalHours", "6", "int", "Intervall für den Parquet/JSON Feedback-Export", DateTime.UtcNow), - new("DefaultLeverageLimit", "10", "decimal", "Standardmäßiger Maximalhebel für Derivate", DateTime.UtcNow) - }; - } + private static readonly ConcurrentDictionary> _inMemorySettings = new(StringComparer.OrdinalIgnoreCase); public AdminSettingsController( WebMqttClient mqttClient, @@ -100,11 +81,57 @@ public class AdminSettingsController : ControllerBase } /// - /// Retrieves all service configurations grouped by service name. + /// Retrieves recent buffered in-memory logs for a specific service. + /// + [HttpGet("logs/{serviceName}")] + public IActionResult GetServiceLogs(string serviceName) + { + if (BackendMqttBridge.ServiceLogsRingBuffer.TryGetValue(serviceName, out var queue)) + { + return Ok(queue.ToList()); + } + return Ok(new List()); + } + + /// + /// Retrieves all service configurations grouped by service name via live MQTT RPC queries. /// [HttpGet] - public IActionResult GetAllSettings() + public async Task GetAllSettings() { + if (_mqttClient.IsConnected) + { + var fetchTasks = ServiceRpcPrefixes.Select(async kvp => + { + var serviceName = kvp.Key; + var prefix = kvp.Value; + try + { + var liveSettings = await _mqttClient.SendRpcRequestAsync, string>( + $"{prefix}_settings_GetAll", + "", + TimeSpan.FromSeconds(2)); + + if (liveSettings != null && liveSettings.Count > 0) + { + _inMemorySettings[serviceName] = liveSettings.Select(d => new ServiceConfigItem( + Key: d.Key, + Value: d.Value?.ToString() ?? "", + DataType: d.Type, + Description: d.Description, + UpdatedAt: d.UpdatedAt ?? DateTime.UtcNow + )).ToList(); + } + } + catch (Exception ex) + { + _logger.LogWarning(ex, "[AdminSettings] Could not fetch live settings from {ServiceName} via MQTT.", serviceName); + } + }); + + await Task.WhenAll(fetchTasks); + } + var grouped = _inMemorySettings.ToDictionary( g => g.Key, g => g.Value.Select(item => new ServiceConfigItemResponseDto( @@ -198,11 +225,37 @@ public class AdminSettingsController : ControllerBase } /// - /// Retrieves settings for a specific service. + /// Retrieves settings for a specific service via live MQTT RPC. /// [HttpGet("{serviceName}")] - public IActionResult GetServiceSettings(string serviceName) + public async Task GetServiceSettings(string serviceName) { + if (ServiceRpcPrefixes.TryGetValue(serviceName, out var prefix) && _mqttClient.IsConnected) + { + try + { + var liveSettings = await _mqttClient.SendRpcRequestAsync, string>( + $"{prefix}_settings_GetAll", + "", + TimeSpan.FromSeconds(2)); + + if (liveSettings != null && liveSettings.Count > 0) + { + _inMemorySettings[serviceName] = liveSettings.Select(d => new ServiceConfigItem( + Key: d.Key, + Value: d.Value?.ToString() ?? "", + DataType: d.Type, + Description: d.Description, + UpdatedAt: d.UpdatedAt ?? DateTime.UtcNow + )).ToList(); + } + } + catch (Exception ex) + { + _logger.LogWarning(ex, "[AdminSettings] Could not fetch live settings for {ServiceName} via MQTT.", serviceName); + } + } + if (_inMemorySettings.TryGetValue(serviceName, out var list)) { var dtos = list.Select(item => new ServiceConfigItemResponseDto( @@ -220,11 +273,10 @@ public class AdminSettingsController : ControllerBase } /// - /// Broadcasts configuration settings to the targeted microservice via MQTT. - /// Does NOT write to Backend database. The microservice persists updated settings directly into its own database. + /// Broadcasts configuration settings to the targeted microservice via MQTT RPC and updates in-memory cache. /// [HttpPut("{serviceName}")] - public async Task UpdateServiceSettings(string serviceName, [FromBody] Dictionary updatedValues) + public async Task UpdateServiceSettings(string serviceName, [FromBody] Dictionary updatedValues) { if (updatedValues == null || !updatedValues.Any()) { @@ -233,40 +285,32 @@ public class AdminSettingsController : ControllerBase _logger.LogInformation("[AdminSettings] Transmitting {Count} config settings to microservice '{ServiceName}' via MQTT", updatedValues.Count, serviceName); - // In-Memory Template-Store aktualisieren - if (_inMemorySettings.TryGetValue(serviceName, out var existingList)) - { - foreach (var (key, value) in updatedValues) - { - var idx = existingList.FindIndex(item => item.Key.Equals(key, StringComparison.OrdinalIgnoreCase)); - if (idx >= 0) - { - var old = existingList[idx]; - existingList[idx] = old with { Value = value, UpdatedAt = DateTime.UtcNow }; - } - else - { - existingList.Add(new ServiceConfigItem(key, value, "string", $"Setting for {serviceName}", DateTime.UtcNow)); - } - } - } - else - { - var newList = updatedValues.Select(kv => new ServiceConfigItem(kv.Key, kv.Value, "string", $"Setting for {serviceName}", DateTime.UtcNow)).ToList(); - _inMemorySettings[serviceName] = newList; - } - - // MQTT Config-Update Event senden bool mqttPublished = false; try { - if (_mqttClient.IsConnected) + if (_mqttClient.IsConnected && ServiceRpcPrefixes.TryGetValue(serviceName, out var prefix)) { + var updated = await _mqttClient.SendRpcRequestAsync, Dictionary>( + $"{prefix}_settings_Update", + updatedValues, + TimeSpan.FromSeconds(3)); + + if (updated != null && updated.Count > 0) + { + _inMemorySettings[serviceName] = updated.Select(d => new ServiceConfigItem( + Key: d.Key, + Value: d.Value?.ToString() ?? "", + DataType: d.Type, + Description: d.Description, + UpdatedAt: d.UpdatedAt ?? DateTime.UtcNow + )).ToList(); + } + string topic = $"services/config/updated/{serviceName}"; var payload = new ServiceConfigUpdatePayload( ServiceName: serviceName, Timestamp: DateTime.UtcNow, - Settings: updatedValues + Settings: updatedValues.ToDictionary(kv => kv.Key, kv => kv.Value?.ToString() ?? "") ); await _mqttClient.PublishAsync(topic, payload); diff --git a/FinlyticBackend/Hubs/LogStreamHub.cs b/FinlyticBackend/Hubs/LogStreamHub.cs new file mode 100644 index 0000000..e0a15b7 --- /dev/null +++ b/FinlyticBackend/Hubs/LogStreamHub.cs @@ -0,0 +1,49 @@ +using System; +using System.Threading.Tasks; +using Microsoft.AspNetCore.SignalR; +using Microsoft.Extensions.Logging; + +namespace FinlyticBackend.Hubs; + +/// +/// SignalR Hub that streams live log messages from microservices to connected web UI clients. +/// +public class LogStreamHub : Hub +{ + private readonly ILogger _logger; + + public LogStreamHub(ILogger logger) + { + _logger = logger; + } + + public async Task JoinServiceLogs(string serviceName) + { + if (!string.IsNullOrWhiteSpace(serviceName)) + { + await Groups.AddToGroupAsync(Context.ConnectionId, serviceName); + _logger.LogInformation("[LogStreamHub] Client {ConnectionId} joined log stream for {ServiceName}", Context.ConnectionId, serviceName); + } + } + + public async Task LeaveServiceLogs(string serviceName) + { + if (!string.IsNullOrWhiteSpace(serviceName)) + { + await Groups.RemoveFromGroupAsync(Context.ConnectionId, serviceName); + _logger.LogInformation("[LogStreamHub] Client {ConnectionId} left log stream for {ServiceName}", Context.ConnectionId, serviceName); + } + } + + public override async Task OnConnectedAsync() + { + _logger.LogInformation("[LogStreamHub] Client connected: {ConnectionId}", Context.ConnectionId); + await base.OnConnectedAsync(); + } + + public override async Task OnDisconnectedAsync(Exception? exception) + { + _logger.LogInformation("[LogStreamHub] Client disconnected: {ConnectionId}", Context.ConnectionId); + await base.OnDisconnectedAsync(exception); + } +} diff --git a/FinlyticBackend/Program.cs b/FinlyticBackend/Program.cs index c733f6f..d8dec91 100644 --- a/FinlyticBackend/Program.cs +++ b/FinlyticBackend/Program.cs @@ -191,6 +191,10 @@ app.MapHub("/hubs/favorites-prices", options => { options.Transports = HttpTransportType.WebSockets | HttpTransportType.ServerSentEvents; }); +app.MapHub("/hubs/logs", options => +{ + options.Transports = HttpTransportType.WebSockets | HttpTransportType.ServerSentEvents; +}); app.MapGet("/health", () => Results.Ok(new { status = "Healthy", service = "FinlyticBackend", timestamp = DateTime.UtcNow })); diff --git a/FinlyticBackend/Util/BackendMqttBridge.cs b/FinlyticBackend/Util/BackendMqttBridge.cs index b1f431e..b6159af 100644 --- a/FinlyticBackend/Util/BackendMqttBridge.cs +++ b/FinlyticBackend/Util/BackendMqttBridge.cs @@ -7,6 +7,7 @@ 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; @@ -28,12 +29,14 @@ 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; @@ -43,6 +46,7 @@ public class BackendMqttBridge : ManagedMqttClient, IHostedService IHubContext hubContext, IHubContext tradeHubContext, IHubContext newsHubContext, + IHubContext logHubContext, IFirebaseNotificationService firebaseService, ILogger logger) : base(logger) { @@ -51,6 +55,7 @@ public class BackendMqttBridge : ManagedMqttClient, IHostedService _hubContext = hubContext; _tradeHubContext = tradeHubContext; _newsHubContext = newsHubContext; + _logHubContext = logHubContext; _firebaseService = firebaseService; _logger = logger; } @@ -95,6 +100,8 @@ public class BackendMqttBridge : ManagedMqttClient, IHostedService await SubscribeAsync("finlytic/technicalanalysis/#"); await SubscribeAsync("finlytic/ta/#"); + // Real-time Logs + await SubscribeAsync("finlytic/logs/#"); } /// @@ -104,7 +111,11 @@ public class BackendMqttBridge : ManagedMqttClient, IHostedService try { - if (topic.StartsWith("finlytic/trades/proposed/", StringComparison.OrdinalIgnoreCase) || + 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); @@ -136,6 +147,23 @@ public class BackendMqttBridge : ManagedMqttClient, IHostedService } } + 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); diff --git a/FinlyticBackend/Util/WebMqttClient.cs b/FinlyticBackend/Util/WebMqttClient.cs index 9d791da..09568f3 100644 --- a/FinlyticBackend/Util/WebMqttClient.cs +++ b/FinlyticBackend/Util/WebMqttClient.cs @@ -61,10 +61,27 @@ public class WebMqttClient : ManagedMqttClient, IHostedService await SubscribeAsync("services/response/assets_GetDiscovery/#"); await SubscribeAsync("services/response/assets_GetDerivatives/#"); await SubscribeAsync("services/response/trades_Get/#"); - await SubscribeAsync("services/response/trades_Close/#"); + await SubscribeAsync("services/response/trades_Reject/#"); + await SubscribeAsync("services/response/trades_Accept/#"); await SubscribeAsync("services/response/analyzer_TriggerManual/#"); await SubscribeAsync("services/response/health_Ping/#"); + + // Settings RPC response channels for all microservices + await SubscribeAsync("services/response/fundamentals_settings_GetAll/#"); + await SubscribeAsync("services/response/fundamentals_settings_Update/#"); + await SubscribeAsync("services/response/news_settings_GetAll/#"); + await SubscribeAsync("services/response/news_settings_Update/#"); + await SubscribeAsync("services/response/ta_settings_GetAll/#"); + await SubscribeAsync("services/response/ta_settings_Update/#"); + await SubscribeAsync("services/response/sentiment_settings_GetAll/#"); + await SubscribeAsync("services/response/sentiment_settings_Update/#"); + await SubscribeAsync("services/response/analyzer_settings_GetAll/#"); + await SubscribeAsync("services/response/analyzer_settings_Update/#"); + await SubscribeAsync("services/response/trades_settings_GetAll/#"); + await SubscribeAsync("services/response/trades_settings_Update/#"); + await SubscribeAsync("services/response/assets_settings_GetAll/#"); + await SubscribeAsync("services/response/assets_settings_Update/#"); } protected override Task OnMessageReceivedAsync(string topic, string payload)