feat(trades): align TradeProposal and TradeAcceptance DTOs, update lifecycle and migrations
This commit is contained in:
@@ -6,6 +6,7 @@ using System.Text.Json;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using FinlyticCore.Dtos;
|
||||
using FinlyticCore.Dtos.TechnicalAnalysis;
|
||||
using FinlyticCore.Models;
|
||||
using FinlyticCore.Models.Trades;
|
||||
using FinlyticCore.Util;
|
||||
@@ -21,18 +22,15 @@ namespace FinlyticTrades.Util;
|
||||
public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
{
|
||||
private readonly IConfiguration _configuration;
|
||||
private readonly ITradeLifecycleService _tradeLifecycleService;
|
||||
private readonly IServiceScopeFactory _scopeFactory;
|
||||
private readonly ILogger<TradesMqttClient> _logger;
|
||||
|
||||
public TradesMqttClient(
|
||||
IConfiguration configuration,
|
||||
ITradeLifecycleService tradeLifecycleService,
|
||||
IServiceScopeFactory scopeFactory,
|
||||
ILogger<TradesMqttClient> logger) : base(logger)
|
||||
{
|
||||
_configuration = configuration;
|
||||
_tradeLifecycleService = tradeLifecycleService;
|
||||
_scopeFactory = scopeFactory;
|
||||
_logger = logger;
|
||||
}
|
||||
@@ -61,7 +59,7 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
protected override async Task OnConnectedAsync()
|
||||
{
|
||||
_logger.LogInformation("[{Channel}] Trades MQTT Client connected. Subscribing to topics...", "TradesChannel");
|
||||
|
||||
|
||||
await SubscribeAsync("finlytic/trades/proposed/#");
|
||||
await SubscribeAsync("finlytic/trades/updates/#");
|
||||
await SubscribeAsync("finlytic/trades/accept/#");
|
||||
@@ -71,6 +69,7 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
await SubscribeAsync("services/request/trades_Accept/#");
|
||||
await SubscribeAsync("services/config/updated/#");
|
||||
await SubscribeAsync("services/request/health_Ping/#");
|
||||
await SubscribeAsync("services/response/tr_GetLivePrice/#");
|
||||
|
||||
_logger.LogInformation("[{Channel}] Successfully subscribed to all event and RPC channels.", "TradesChannel");
|
||||
}
|
||||
@@ -114,12 +113,16 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
return;
|
||||
}
|
||||
|
||||
// Für Scoped-Services erzeugen wir pro eingehender Nachricht einen eigenen Scope
|
||||
using var msgScope = _scopeFactory.CreateScope();
|
||||
var tradeLifecycleService = msgScope.ServiceProvider.GetRequiredService<ITradeLifecycleService>();
|
||||
|
||||
if (topic.StartsWith("finlytic/trades/proposed/"))
|
||||
{
|
||||
var proposal = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.TradeProposalDto);
|
||||
if (proposal != null && (!string.IsNullOrWhiteSpace(proposal.Symbol) || !string.IsNullOrWhiteSpace(proposal.Isin)))
|
||||
{
|
||||
await _tradeLifecycleService.ProcessProposedTradeAsync(proposal, CancellationToken.None);
|
||||
await tradeLifecycleService.ProcessProposedTradeAsync(proposal, CancellationToken.None);
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -131,7 +134,7 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
var acceptDto = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.TradeAcceptanceDto);
|
||||
if (acceptDto != null)
|
||||
{
|
||||
var newTrade = await _tradeLifecycleService.AcceptTradeAsync(acceptDto, CancellationToken.None);
|
||||
var newTrade = await tradeLifecycleService.AcceptTradeAsync(acceptDto, CancellationToken.None);
|
||||
if (newTrade != null)
|
||||
{
|
||||
var dto = MapToDto(newTrade);
|
||||
@@ -145,7 +148,7 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
var acceptDto = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.TradeAcceptanceDto);
|
||||
if (acceptDto != null)
|
||||
{
|
||||
var acceptedTrade = await _tradeLifecycleService.AcceptTradeAsync(acceptDto, CancellationToken.None);
|
||||
var acceptedTrade = await tradeLifecycleService.AcceptTradeAsync(acceptDto, CancellationToken.None);
|
||||
if (acceptedTrade != null)
|
||||
{
|
||||
var acceptedDto = MapToDto(acceptedTrade);
|
||||
@@ -159,21 +162,48 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
var update = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.TradeHourlyUpdateDto);
|
||||
if (update != null)
|
||||
{
|
||||
await _tradeLifecycleService.AddHourlyUpdateAsync(update, CancellationToken.None);
|
||||
await tradeLifecycleService.AddHourlyUpdateAsync(update, CancellationToken.None);
|
||||
}
|
||||
}
|
||||
else if (topic.StartsWith("services/request/trades_Get/"))
|
||||
{
|
||||
var correlationId = topic.Split('/').Last();
|
||||
var request = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.GetTradesRequest);
|
||||
|
||||
|
||||
string? isin = request?.Isin;
|
||||
string? status = request?.Status;
|
||||
string? userId = request?.UserId;
|
||||
|
||||
var trades = await _tradeLifecycleService.GetTradesAsync(isin, status, userId);
|
||||
var trades = await tradeLifecycleService.GetTradesAsync(isin, status, userId);
|
||||
|
||||
var activeTrades = trades.Where(t => t.Status == TradeStatus.Active && !string.IsNullOrWhiteSpace(t.Isin)).ToList();
|
||||
if (activeTrades.Count > 0)
|
||||
{
|
||||
try
|
||||
{
|
||||
var priceTasks = activeTrades.Select(t => FetchLivePriceAsync(t.Isin)).ToList();
|
||||
var livePricesTask = Task.WhenAll(priceTasks);
|
||||
if (await Task.WhenAny(livePricesTask, Task.Delay(1500)) == livePricesTask)
|
||||
{
|
||||
var livePrices = await livePricesTask;
|
||||
for (int i = 0; i < activeTrades.Count; i++)
|
||||
{
|
||||
var lp = livePrices[i];
|
||||
if (lp != null && lp.CurrentPrice > 0m)
|
||||
{
|
||||
var trade = activeTrades[i];
|
||||
tradeLifecycleService.CalculatePnL(trade, lp.CurrentPrice);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogDebug(ex, "[{Channel}] Live price fetch skipped or timed out during trades_Get", "TradesChannel");
|
||||
}
|
||||
}
|
||||
|
||||
var dtos = trades.Select(MapToDto).ToList();
|
||||
|
||||
await PublishAsync($"services/response/trades_Get/{correlationId}", dtos);
|
||||
}
|
||||
else if (topic.StartsWith("services/request/trades_Close/"))
|
||||
@@ -181,18 +211,17 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
var parts = topic.Split('/');
|
||||
var tradeId = parts.Length > 3 ? parts[3] : string.Empty;
|
||||
var correlationId = parts.Length > 4 ? parts[4] : string.Empty;
|
||||
|
||||
|
||||
var request = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.CloseTradeRequest);
|
||||
|
||||
if (request != null && !string.IsNullOrEmpty(tradeId))
|
||||
{
|
||||
var closedTrade = await _tradeLifecycleService.CloseTradeAsync(tradeId, request);
|
||||
var closedTrade = await tradeLifecycleService.CloseTradeAsync(tradeId, request);
|
||||
if (closedTrade != null)
|
||||
{
|
||||
var closedDto = MapToDto(closedTrade);
|
||||
await PublishAsync($"services/response/trades_Close/{correlationId}", closedDto);
|
||||
|
||||
// Send event stream update specifically for closed trades (used by Feedback Engine & Analytics)
|
||||
|
||||
string sectorSafe = string.IsNullOrWhiteSpace(closedTrade.Sector) ? "general" : closedTrade.Sector.ToLowerInvariant();
|
||||
await PublishAsync($"finlytic/trades/closed/{sectorSafe}/{closedTrade.Symbol.ToLowerInvariant()}", closedDto);
|
||||
await PublishTradeUpdateAsync(closedDto);
|
||||
@@ -204,12 +233,12 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
var parts = topic.Split('/');
|
||||
var tradeId = parts.Length > 3 ? parts[3] : string.Empty;
|
||||
var correlationId = parts.Length > 4 ? parts[4] : string.Empty;
|
||||
|
||||
|
||||
var request = JsonSerializer.Deserialize(payloadStr, FinlyticJsonSerializerContext.Default.CloseTradeRequest);
|
||||
|
||||
if (request != null && !string.IsNullOrEmpty(tradeId))
|
||||
{
|
||||
var rejectedTrade = await _tradeLifecycleService.RejectTradeAsync(tradeId, request);
|
||||
var rejectedTrade = await tradeLifecycleService.RejectTradeAsync(tradeId, request);
|
||||
if (rejectedTrade != null)
|
||||
{
|
||||
var rejectedDto = MapToDto(rejectedTrade);
|
||||
@@ -227,9 +256,26 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
|
||||
public async Task PublishTradeUpdateAsync(TradeProposalDto trade)
|
||||
{
|
||||
await PublishAsync($"finlytic/trades/user/{trade.UserId ?? "all"}", trade);
|
||||
await PublishAsync("finlytic/trades/update", trade);
|
||||
}
|
||||
|
||||
private async Task<LivePriceDto?> FetchLivePriceAsync(string isin)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(isin)) return null;
|
||||
try
|
||||
{
|
||||
return await SendRpcRequestAsync<LivePriceDto, IsinRequest>(
|
||||
"tr_GetLivePrice",
|
||||
new IsinRequest(isin),
|
||||
TimeSpan.FromMilliseconds(1200));
|
||||
}
|
||||
catch
|
||||
{
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private static TradeProposalDto MapToDto(TradeEntity t)
|
||||
{
|
||||
List<decimal>? parseTakeProfitTargets()
|
||||
@@ -279,7 +325,7 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
FundamentalRationale = t.FundamentalRationale,
|
||||
RiskWarning = t.RiskWarning,
|
||||
CreatedAt = t.CreatedAt,
|
||||
|
||||
|
||||
UserId = t.UserId,
|
||||
IsGlobalProposal = t.IsGlobalProposal,
|
||||
ActualEntryPrice = t.ActualEntryPrice,
|
||||
@@ -290,7 +336,10 @@ public class TradesMqttClient : ManagedMqttClient, IHostedService
|
||||
ExecutionTimestamp = t.ExecutionTimestamp,
|
||||
Quantity = t.Quantity,
|
||||
KnockoutThreshold = t.KnockoutThreshold,
|
||||
IsRecurring = t.IsRecurring
|
||||
IsRecurring = t.IsRecurring,
|
||||
PnlAbsolute = t.PnlAbsolute,
|
||||
PnlPercent = t.PnlPercent,
|
||||
CurrentPrice = t.UserExitPrice ?? t.HourlyUpdates?.LastOrDefault()?.CurrentPrice
|
||||
};
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user