using System.Text.Json; using KArtSell.BuildingBlocks.Data; using KArtSell.BuildingBlocks.Hashing; using KArtSell.BuildingBlocks.Reliability; using KArtSell.BuildingBlocks.Time; using Microsoft.Extensions.Logging; namespace KArtSell.Modules.ModelOperations.TradeExecution; public class SubmitTradeCommand { public Guid SellDecisionId { get; set; } public int Quantity { get; set; } public decimal LimitPrice { get; set; } public Guid CorrelationId { get; set; } } public class SubmitTradeHandler { private readonly ITradeSql _sql; private readonly IKisTradeExecutionService _kis; private readonly IDbConnectionFactory _connectionFactory; private readonly IOutboxWriter _outbox; private readonly IClock _clock; private readonly ILogger _logger; public SubmitTradeHandler( ITradeSql sql, IKisTradeExecutionService kis, IDbConnectionFactory connectionFactory, IOutboxWriter outbox, IClock clock, ILogger logger) { _sql = sql; _kis = kis; _connectionFactory = connectionFactory; _outbox = outbox; _clock = clock; _logger = logger; } public async Task HandleAsync(SubmitTradeCommand command, CancellationToken ct = default) { var trade = Trade.Create(command.SellDecisionId, command.Quantity, command.CorrelationId, _clock.UtcNow.UtcDateTime); await _sql.InsertTradeAsync(trade, ct); _logger.LogInformation("Created trade: {TradeId}", trade.Id); try { var (orderId, response) = await _kis.ExecuteTradeAsync( trade.Id, command.Quantity, command.LimitPrice, command.CorrelationId, ct ); trade.MarkSubmitted(orderId, response); await _sql.UpdateTradeStatusAsync(trade, response, null, ct); await PublishEventAsync( "TradeSubmitted", new TradeSubmittedEvent { TradeId = trade.Id, SellDecisionId = command.SellDecisionId, KisOrderId = orderId, Quantity = command.Quantity, CorrelationId = command.CorrelationId }, command.CorrelationId, ct); return trade.Id; } catch (KisTradeExecutionException ex) { trade.MarkErrored(ex); await _sql.UpdateTradeStatusAsync(trade, ex.KisResponse, ex.Message, ct); _logger.LogError( "Trade submission failed: {TradeId} {Classification}", trade.Id, ex.Classification ); throw; } } private async Task PublishEventAsync(string eventType, T @event, Guid correlationId, CancellationToken ct) where T : class => await TradeOutboxPublisher.PublishAsync(_connectionFactory, _outbox, _clock, eventType, @event, correlationId, ct); } public class PollTradeStatusCommand { public Guid TradeId { get; set; } public string KisOrderId { get; set; } = string.Empty; public Guid CorrelationId { get; set; } } public class PollTradeStatusHandler { private readonly ITradeSql _sql; private readonly IKisTradeExecutionService _kis; private readonly IDbConnectionFactory _connectionFactory; private readonly IOutboxWriter _outbox; private readonly IClock _clock; private readonly ILogger _logger; public PollTradeStatusHandler( ITradeSql sql, IKisTradeExecutionService kis, IDbConnectionFactory connectionFactory, IOutboxWriter outbox, IClock clock, ILogger logger) { _sql = sql; _kis = kis; _connectionFactory = connectionFactory; _outbox = outbox; _clock = clock; _logger = logger; } public async Task HandleAsync(PollTradeStatusCommand command, CancellationToken ct = default) { var trade = await _sql.GetTradeByIdAsync(command.TradeId, command.CorrelationId, ct); if (trade == null) { _logger.LogWarning("Trade not found: {TradeId}", command.TradeId); return; } try { var (status, executedQty, unitPrice, response) = await _kis.GetOrderStatusAsync( command.KisOrderId, command.CorrelationId, ct ); if (status is "ACCEPTED" or "PARTIAL_FILLED" or "FULLY_FILLED") { trade.MarkAccepted(response); if (status is "PARTIAL_FILLED" or "FULLY_FILLED") { trade.MarkFilled(executedQty, unitPrice, response, _clock.UtcNow.UtcDateTime); } await _sql.UpdateTradeStatusAsync(trade, response, null, ct); if (trade.Status is TradeStatus.FullyFilled) { await TradeOutboxPublisher.PublishAsync( _connectionFactory, _outbox, _clock, "TradeFilled", new TradeFilledEvent { TradeId = trade.Id, ExecutedQuantity = executedQty, UnitPrice = unitPrice, CorrelationId = command.CorrelationId }, command.CorrelationId, ct); } _logger.LogInformation("Trade status updated: {TradeId} -> {Status}", trade.Id, status); } } catch (KisTradeExecutionException ex) { await _sql.UpdateTradeStatusAsync(trade, ex.KisResponse, ex.Message, ct); _logger.LogError("Failed to poll trade status: {TradeId}", trade.Id); } } } public class ConfirmSettlementCommand { public Guid TradeId { get; set; } public string KisOrderId { get; set; } = string.Empty; public decimal? Commission { get; set; } public Guid CorrelationId { get; set; } } public class ConfirmSettlementHandler { private readonly ITradeSql _sql; private readonly IKisTradeExecutionService _kis; private readonly IDbConnectionFactory _connectionFactory; private readonly IOutboxWriter _outbox; private readonly IClock _clock; private readonly ILogger _logger; public ConfirmSettlementHandler( ITradeSql sql, IKisTradeExecutionService kis, IDbConnectionFactory connectionFactory, IOutboxWriter outbox, IClock clock, ILogger logger) { _sql = sql; _kis = kis; _connectionFactory = connectionFactory; _outbox = outbox; _clock = clock; _logger = logger; } public async Task HandleAsync(ConfirmSettlementCommand command, CancellationToken ct = default) { var trade = await _sql.GetTradeByIdAsync(command.TradeId, command.CorrelationId, ct); if (trade == null) { _logger.LogWarning("Trade not found for settlement: {TradeId}", command.TradeId); return; } try { var (success, response) = await _kis.ConfirmSettlementAsync( command.KisOrderId, command.CorrelationId, ct ); if (success) { trade.MarkConfirmed(_clock.UtcNow.UtcDateTime, command.Commission); await _sql.UpdateTradeStatusAsync(trade, response, null, ct); await TradeOutboxPublisher.PublishAsync( _connectionFactory, _outbox, _clock, "TradeSettled", new TradeSettledEvent { TradeId = trade.Id, NetProceeds = trade.NetProceeds ?? 0, CorrelationId = command.CorrelationId }, command.CorrelationId, ct); _logger.LogInformation("Trade settlement confirmed: {TradeId}", trade.Id); } } catch (KisTradeExecutionException ex) { await _sql.UpdateTradeStatusAsync(trade, ex.KisResponse, ex.Message, ct); _logger.LogError("Failed to confirm settlement: {TradeId}", trade.Id); } } } public class TradeSubmittedEvent { public Guid TradeId { get; set; } public Guid SellDecisionId { get; set; } public string KisOrderId { get; set; } = string.Empty; public int Quantity { get; set; } public Guid CorrelationId { get; set; } } public class TradeFilledEvent { public Guid TradeId { get; set; } public int ExecutedQuantity { get; set; } public decimal UnitPrice { get; set; } public Guid CorrelationId { get; set; } } public class TradeSettledEvent { public Guid TradeId { get; set; } public decimal NetProceeds { get; set; } public Guid CorrelationId { get; set; } } /// /// DEBT-TRADE-001: outbox write happens in its own transaction, separate from the /// preceding trade status update (which owns its own connection in TradeSql). Not yet /// atomic with the state transition. See TECH_DEBT_REGISTER.md. /// internal static class TradeOutboxPublisher { public static async Task PublishAsync( IDbConnectionFactory connectionFactory, IOutboxWriter outbox, IClock clock, string eventType, T @event, Guid correlationId, CancellationToken ct) where T : class { var payload = JsonSerializer.Serialize(@event); var message = new OutboxMessage( Guid.NewGuid(), eventType, 1, payload, correlationId.ToString(), clock.UtcNow, ContentHasher.Sha256(payload)); await using var connection = await connectionFactory.OpenAsync(ct); await using var transaction = await connection.BeginTransactionAsync(ct); await outbox.AddAsync(connection, transaction, message, ct); await transaction.CommitAsync(ct); } }