feat(Backend): update API gateway and websocket hubs

This commit is contained in:
2026-08-09 21:32:52 +02:00
parent fdf4b6efcb
commit a9553e9fbf
66 changed files with 188878 additions and 0 deletions
@@ -0,0 +1,22 @@
using Microsoft.AspNetCore.Mvc.Filters;
namespace FinlyticBackend.Util;
/// <summary>
/// A global authorization filter that short-circuits authorization checks for HTTP OPTIONS requests.
/// This is required for CORS preflight to succeed on endpoints protected with [Authorize]:
/// the browser sends a parameter-less OPTIONS request before the real request, and any 401
/// response on that preflight causes the actual request to be blocked with a CORS error.
/// </summary>
public class AllowOptionsFilter : IAuthorizationFilter
{
public void OnAuthorization(AuthorizationFilterContext context)
{
if (context.HttpContext.Request.Method == HttpMethods.Options)
{
// Return 204 No Content immediately — the manual CORS middleware in Program.cs
// has already written the Access-Control-Allow-* headers.
context.Result = new Microsoft.AspNetCore.Mvc.StatusCodeResult(204);
}
}
}
+220
View File
@@ -0,0 +1,220 @@
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using FinlyticBackend.Database;
using FinlyticBackend.Hubs;
using FinlyticBackend.Services;
using FinlyticCore.Dtos.News;
using FinlyticCore.Models;
using FinlyticCore.Models.Auth;
using FinlyticCore.Models.Trades;
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;
/// <summary>
/// Subscribes to general broadcast MQTT topics and forwards them to SignalR clients and FCM push services.
/// </summary>
public class BackendMqttBridge : ManagedMqttClient, IHostedService
{
public static readonly ConcurrentDictionary<string, JsonElement> FundamentalsCache = new(StringComparer.OrdinalIgnoreCase);
public static readonly ConcurrentDictionary<string, JsonElement> TechnicalsCache = new(StringComparer.OrdinalIgnoreCase);
private readonly IConfiguration _configuration;
private readonly IServiceScopeFactory _scopeFactory;
private readonly IHubContext<TradeRealtimeHub, ITradeClient> _hubContext;
private readonly IHubContext<TradeHub> _tradeHubContext;
private readonly IHubContext<NewsHub> _newsHubContext;
private readonly IFirebaseNotificationService _firebaseService;
private readonly ILogger<BackendMqttBridge> _logger;
public BackendMqttBridge(
IConfiguration configuration,
IServiceScopeFactory scopeFactory,
IHubContext<TradeRealtimeHub, ITradeClient> hubContext,
IHubContext<TradeHub> tradeHubContext,
IHubContext<NewsHub> newsHubContext,
IFirebaseNotificationService firebaseService,
ILogger<BackendMqttBridge> logger) : base(logger)
{
_configuration = configuration;
_scopeFactory = scopeFactory;
_hubContext = hubContext;
_tradeHubContext = tradeHubContext;
_newsHubContext = newsHubContext;
_firebaseService = firebaseService;
_logger = logger;
}
/// <inheritdoc />
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_bridge")}_{Guid.NewGuid()}"
};
_logger.LogInformation("Starting Backend MQTT Bridge. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId);
await ConnectAsync(config);
}
/// <inheritdoc />
public async Task StopAsync(CancellationToken cancellationToken)
{
_logger.LogInformation("Stopping Backend MQTT Bridge.");
await DisconnectAsync();
}
/// <inheritdoc />
protected override async Task OnConnectedAsync()
{
_logger.LogInformation("Backend MQTT Bridge connected. Subscribing to broadcast topics...");
await SubscribeAsync("finlytic/trades/proposed/#");
await SubscribeAsync("finlytic/trades/updates/#");
// News topics
await SubscribeAsync("services/news/completed");
await SubscribeAsync("finlytic/news/#");
await SubscribeAsync("finlytic/sentiment/#");
// Fundamentals & Technicals
await SubscribeAsync("finlytic/fundamentals/#");
await SubscribeAsync("finlytic/assets/fundamentals/#");
await SubscribeAsync("finlytic/technicalanalysis/#");
await SubscribeAsync("finlytic/ta/#");
}
/// <inheritdoc />
protected override async Task OnMessageReceivedAsync(string topic, string payloadStr)
{
if (string.IsNullOrWhiteSpace(topic) || string.IsNullOrWhiteSpace(payloadStr)) return;
try
{
if (topic.StartsWith("finlytic/trades/proposed/", StringComparison.OrdinalIgnoreCase) ||
topic.StartsWith("finlytic/trades/update", StringComparison.OrdinalIgnoreCase))
{
await HandleTradeProposalAsync(payloadStr);
}
else if (topic.StartsWith("finlytic/trades/updates/", StringComparison.OrdinalIgnoreCase))
{
await HandleTradeUpdateAsync(payloadStr);
}
else if (topic.Equals("services/news/completed", StringComparison.OrdinalIgnoreCase) ||
topic.StartsWith("finlytic/news/", StringComparison.OrdinalIgnoreCase) ||
topic.StartsWith("finlytic/sentiment/", StringComparison.OrdinalIgnoreCase))
{
await HandleNewsArticleAsync(payloadStr);
}
else if (topic.StartsWith("finlytic/fundamentals/", StringComparison.OrdinalIgnoreCase) ||
topic.StartsWith("finlytic/assets/fundamentals/", StringComparison.OrdinalIgnoreCase))
{
HandleFundamentalsCache(topic, payloadStr);
}
else if (topic.StartsWith("finlytic/technicalanalysis/", StringComparison.OrdinalIgnoreCase) ||
topic.StartsWith("finlytic/ta/", StringComparison.OrdinalIgnoreCase))
{
HandleTechnicalsCache(topic, payloadStr);
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing message in Backend MQTT Bridge on topic {Topic}", topic);
}
}
private async Task HandleTradeProposalAsync(string payloadStr)
{
var proposal = JsonSerializer.Deserialize<TradeProposalDto>(payloadStr);
if (proposal == null) return;
await _hubContext.Clients.All.OnTradeProposed(proposal);
await _tradeHubContext.Clients.All.SendAsync("ReceiveTradeUpdate", proposal);
_logger.LogInformation("Broadcasted Trade Proposal {AnalysisId} via SignalR.", proposal.AnalysisId);
using var scope = _scopeFactory.CreateScope();
var dbContext = scope.ServiceProvider.GetRequiredService<BackendDbContext>();
var fcmTokens = await dbContext.UserDeviceTokens.AsNoTracking().Select(t => t.FcmToken).ToListAsync();
if (fcmTokens.Count > 0)
{
await _firebaseService.SendTradeProposalNotificationAsync(proposal, fcmTokens);
}
}
private async Task HandleTradeUpdateAsync(string payloadStr)
{
var update = JsonSerializer.Deserialize<TradeHourlyUpdateDto>(payloadStr);
if (update == null) return;
await _hubContext.Clients.All.OnTradeUpdated(update);
await _tradeHubContext.Clients.All.SendAsync("ReceiveTradeUpdate", update);
if (string.Equals(update.Recommendation, "Close", StringComparison.OrdinalIgnoreCase) ||
string.Equals(update.Recommendation, "AdjustSL", StringComparison.OrdinalIgnoreCase))
{
using var scope = _scopeFactory.CreateScope();
var dbContext = scope.ServiceProvider.GetRequiredService<BackendDbContext>();
var fcmTokens = await dbContext.UserDeviceTokens.AsNoTracking().Select(t => t.FcmToken).ToListAsync();
if (fcmTokens.Count > 0)
{
await _firebaseService.SendTradeUpdateNotificationAsync(update, fcmTokens);
}
}
}
private async Task HandleNewsArticleAsync(string payloadStr)
{
var article = JsonSerializer.Deserialize<NewsArticleDto>(payloadStr, FinlyticJsonSerializerContext.Default.NewsArticleDto);
if (article == null) return;
await _newsHubContext.Clients.All.SendAsync("ReceiveNewArticle", article);
_logger.LogInformation("Broadcasted live news item '{Title}' over SignalR NewsHub.", article.Title);
}
private void HandleFundamentalsCache(string topic, string payloadStr)
{
using var jsonDoc = JsonDocument.Parse(payloadStr);
var root = jsonDoc.RootElement.Clone();
string? isin = root.TryGetProperty("isin", out var isinProp) ? isinProp.GetString() : null;
string? ticker = root.TryGetProperty("primaryTicker", out var tickerProp)
? tickerProp.GetString()
: (root.TryGetProperty("ticker", out var tProp) ? tProp.GetString() : null);
if (!string.IsNullOrEmpty(isin)) FundamentalsCache[isin] = root;
if (!string.IsNullOrEmpty(ticker)) FundamentalsCache[ticker] = root;
string topicKey = topic.Split('/').LastOrDefault() ?? string.Empty;
if (!string.IsNullOrEmpty(topicKey)) FundamentalsCache[topicKey] = root;
_logger.LogInformation("Cached fundamentals payload from MQTT topic {Topic}.", topic);
}
private void HandleTechnicalsCache(string topic, string payloadStr)
{
using var jsonDoc = JsonDocument.Parse(payloadStr);
var root = jsonDoc.RootElement.Clone();
string? symbol = root.TryGetProperty("symbol", out var sProp) ? sProp.GetString() : null;
if (!string.IsNullOrEmpty(symbol)) TechnicalsCache[symbol] = root;
string topicKey = topic.Split('/').LastOrDefault() ?? string.Empty;
if (!string.IsNullOrEmpty(topicKey)) TechnicalsCache[topicKey] = root;
_logger.LogInformation("Cached technical analysis payload from MQTT topic {Topic}.", topic);
}
}
+11
View File
@@ -0,0 +1,11 @@
namespace FinlyticAssets.Util;
public class Volumes
{
/// <summary>
/// Der relative Pfad für die schlanke Index-Datei (ISINs + Namen) zur Asset-Erkennung.
/// </summary>
public const string IndexRelativePath = "assets/index";
public const string LogosRelativePath = "assets/logos";
}
+70
View File
@@ -0,0 +1,70 @@
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;
/// <summary>
/// Managed MQTT client for web gateway endpoints enabling RPC communication with background microservices.
/// </summary>
public class WebMqttClient : ManagedMqttClient, IHostedService
{
private readonly ILogger<WebMqttClient> _logger;
private readonly IConfiguration _configuration;
public WebMqttClient(ILogger<WebMqttClient> 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/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;
}
}