using System.Threading.Channels; using MongoDB.Driver; using PolyTraderSharp.Extensions; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services { public class PersistenceService : BackgroundService { private readonly ChannelReader _tradeReader; private readonly IMongoDatabase _db; private readonly TerminalLogger _logger; private readonly JobStatusRow _jobStatus; public PersistenceService(ChannelReader tradeReader, IMongoDatabase db, TerminalLogger logger, JobManager jobManager) { _tradeReader = tradeReader; _db = db; _logger = logger; _jobStatus = new JobStatusRow { JobName = "MongoDB 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) var col = _db.GetCollection("closed_trades"); col.EnsureIndex(x => x.TradeId); col.EnsureIndex(x => x.AccountId); col.EnsureIndex(x => x.TokenId); // 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)) { var existing = col.LiteFindOne(x => x.AccountId == trade.AccountId && x.TokenId == trade.TokenId); if (existing != null) { _logger.Debug($"Duplikat ignoriert: ClosedTrade für Account {trade.AccountId} + Token {trade.TokenId.Substring(0, Math.Min(10, trade.TokenId.Length))}... existiert bereits (DB-ID: {existing.TradeId})."); continue; } } col.Insert(trade); _logger.Debug($"Saved ClosedTrade {trade.TradeId} to MongoDB"); _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)..."; } } } } }