using System.Threading.Channels; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using PolyTrader.Core.Persistence; using PolyTrader.Modules.CopyTrading.Persistence; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services { public class PersistenceService : BackgroundService { private readonly ChannelReader _tradeReader; private readonly ICopyTradeLogRepository _tradeLog; private readonly ITradeLogRepository _coreTradeLog; private readonly TerminalLogger _logger; private readonly JobStatusRow _jobStatus; public PersistenceService(ChannelReader tradeReader, ICopyTradeLogRepository tradeLog, ITradeLogRepository coreTradeLog, TerminalLogger logger, JobManager jobManager) { _tradeReader = tradeReader; _tradeLog = tradeLog; _coreTradeLog = coreTradeLog; _logger = logger; _jobStatus = new JobStatusRow { JobName = "MySQL Transaction Log", Description = "Awaits internal signals to write Closed Trades to the database safely.", StatusText = "Pending Initial Delay..." }; _jobStatus.ManualTriggerAction = async () => { _jobStatus.StatusText = "Manual trigger not supported for Channel Reader"; await Task.Delay(2000); _jobStatus.StatusText = "Listening (Channel)..."; }; jobManager.RegisterJob(_jobStatus); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.Info("PersistenceService started writing background DB logs."); _jobStatus.StatusText = "Listening (Channel)..."; // One-time index setup (moved out of hot loop) _tradeLog.EnsureIndexes(); _coreTradeLog.EnsureIndexes(); // We do a loop waiting for items in the channel await foreach(var trade in _tradeReader.ReadAllAsync(stoppingToken)) { if (!_jobStatus.IsEnabled) { // If paused, we just drop the trade for now or log a warning _logger.Warning("PersistenceService is paused, ignoring trade log."); continue; } try { _jobStatus.StatusText = "Writing to DB..."; // ===== DEDUPLIZIERUNG: Verhindert das Mehrfach-Einfügen desselben Trades ===== // Prüft ob für diesen Account + TokenId bereits ein ClosedTrade existiert. // Dies verhindert den "Background Sync Duplicate Bug", bei dem geschlossene // Trades bei jedem Sync-Zyklus oder nach einem Neustart erneut eingefügt werden. if (!string.IsNullOrEmpty(trade.TokenId)) { if (_tradeLog.Exists(trade.AccountId, trade.TokenId)) { _logger.Debug($"Duplikat ignoriert: ClosedTrade für Account {trade.AccountId} + Token {trade.TokenId.Substring(0, Math.Min(10, trade.TokenId.Length))}... existiert bereits."); continue; } } _tradeLog.Insert(trade); // Dual-Write: generischer, modulübergreifender Core-Trade-Log (fürs Dashboard). _coreTradeLog.Insert(new TradeRecord { ModuleName = "CopyTrading", AccountId = trade.AccountId, IsDemo = trade.IsDemo, TokenId = trade.TokenId, MarketQuestion = trade.MarketQuestion, Outcome = trade.Outcome, Side = trade.Side, EntryPrice = trade.EntryPrice, ExitPrice = trade.ExitPrice, Size = trade.Size, RealizedPnl = trade.RealizedPnl, PnlPercent = trade.PnlPercent, OpenedAt = trade.OpenedAt, ClosedAt = trade.ClosedAt, ExitReason = trade.ExitReason }); _logger.Debug($"Saved ClosedTrade {trade.TradeId} to MySQL"); _jobStatus.LastRun = DateTime.Now; } catch (Exception ex) { _logger.Error($"Failed to persist ClosedTrade (ID: {trade.TradeId}): {ex.Message}"); _jobStatus.StatusText = "Error!"; } finally { if (_jobStatus.StatusText != "Error!") _jobStatus.StatusText = "Listening (Channel)..."; } } } } }