249 lines
11 KiB
C#
249 lines
11 KiB
C#
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.Logging;
|
|
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);
|
|
public static readonly ConcurrentDictionary<string, ConcurrentQueue<LogMessageDto>> ServiceLogsRingBuffer = 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 IHubContext<LogStreamHub> _logHubContext;
|
|
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,
|
|
IHubContext<LogStreamHub> logHubContext,
|
|
IFirebaseNotificationService firebaseService,
|
|
ILogger<BackendMqttBridge> logger) : base(logger)
|
|
{
|
|
_configuration = configuration;
|
|
_scopeFactory = scopeFactory;
|
|
_hubContext = hubContext;
|
|
_tradeHubContext = tradeHubContext;
|
|
_newsHubContext = newsHubContext;
|
|
_logHubContext = logHubContext;
|
|
_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/#");
|
|
|
|
// Real-time Logs
|
|
await SubscribeAsync("finlytic/logs/#");
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
protected override async Task OnMessageReceivedAsync(string topic, string payloadStr)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(topic) || string.IsNullOrWhiteSpace(payloadStr)) return;
|
|
|
|
try
|
|
{
|
|
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);
|
|
}
|
|
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 HandleLogMessageAsync(string payloadStr)
|
|
{
|
|
var logDto = JsonSerializer.Deserialize<LogMessageDto>(payloadStr, FinlyticJsonSerializerContext.Default.LogMessageDto);
|
|
if (logDto == null) return;
|
|
|
|
string serviceKey = logDto.ServiceName;
|
|
var queue = ServiceLogsRingBuffer.GetOrAdd(serviceKey, _ => new ConcurrentQueue<LogMessageDto>());
|
|
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<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);
|
|
}
|
|
} |