From 88fc982089e5c79725d26b15ad9e279edd65d09f Mon Sep 17 00:00:00 2001 From: Richard Date: Wed, 1 Jul 2026 22:27:22 +0200 Subject: [PATCH] Phase 5.3-ClosedTrade: ClosedTrade -> Modul + ICopyTradeLogRepository - ClosedTrade (+ ClosedTradeRow) nach src/PolyTrader.Modules.CopyTrading/Models/. - Neues ICopyTradeLogRepository (+ Mongo-Impl) im Modul: EnsureIndexes, Exists, Insert, Find(predicate) auf der closed_trades-Collection. - closed_trades-Zugriffe der Modul-Services vom App-Shim auf das Repo umgestellt: TraderMonitorService, PersistenceService, PolymarketWssClient, TraderAnalyticsJob. Shim-/DB-Usings dort entfernt. - TraderMonitorService nutzt kein _db mehr (auch die uebrigen _db-Guards aus 3d entfernt, da Repos immer verfuegbar sind). Verhaltensneutral. - frm_main (UI) + Program.cs (BsonDocument-Cleanup) bleiben auf _db. - Build 0 Fehler. Co-Authored-By: Claude Opus 4.8 --- Program.cs | 2 + services/PersistenceService.cs | 23 +++---- services/PolymarketWssClient.cs | 11 ++- services/TraderAnalyticsJob.cs | 20 +++--- services/TraderMonitorService.cs | 69 ++++++++----------- .../Models}/ClosedTrade.cs | 0 .../Persistence/ICopyTradeLogRepository.cs | 25 +++++++ .../MongoCopyTradeLogRepository.cs | 45 ++++++++++++ 8 files changed, 122 insertions(+), 73 deletions(-) rename {Models => src/PolyTrader.Modules.CopyTrading/Models}/ClosedTrade.cs (100%) create mode 100644 src/PolyTrader.Modules.CopyTrading/Persistence/ICopyTradeLogRepository.cs create mode 100644 src/PolyTrader.Modules.CopyTrading/Persistence/MongoCopyTradeLogRepository.cs diff --git a/Program.cs b/Program.cs index 8adbf43..4a3cd94 100644 --- a/Program.cs +++ b/Program.cs @@ -10,6 +10,7 @@ using Microsoft.Extensions.Options; using PolyTrader.Core.Configuration; using PolyTrader.Core.DependencyInjection; using PolyTrader.Core.Streaming; +using PolyTrader.Modules.CopyTrading.Persistence; using PolyTraderSharp.Models; using PolyTraderSharp.Services; @@ -36,6 +37,7 @@ internal static class Program return client.GetDatabase(dbOptions.DatabaseName); }); services.AddCorePersistence(); + services.AddSingleton(); services.AddSingleton(); services.AddSingleton((IServiceProvider sp) => ServerSettings.Load("server_settings.xml")); services.AddSingleton(); diff --git a/services/PersistenceService.cs b/services/PersistenceService.cs index 90db44b..4230435 100644 --- a/services/PersistenceService.cs +++ b/services/PersistenceService.cs @@ -1,8 +1,7 @@ using System.Threading.Channels; -using MongoDB.Driver; -using PolyTraderSharp.Extensions; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using PolyTrader.Modules.CopyTrading.Persistence; using PolyTraderSharp.Models; namespace PolyTraderSharp.Services @@ -10,14 +9,14 @@ namespace PolyTraderSharp.Services public class PersistenceService : BackgroundService { private readonly ChannelReader _tradeReader; - private readonly IMongoDatabase _db; + private readonly ICopyTradeLogRepository _tradeLog; private readonly TerminalLogger _logger; private readonly JobStatusRow _jobStatus; - public PersistenceService(ChannelReader tradeReader, IMongoDatabase db, TerminalLogger logger, JobManager jobManager) + public PersistenceService(ChannelReader tradeReader, ICopyTradeLogRepository tradeLog, TerminalLogger logger, JobManager jobManager) { _tradeReader = tradeReader; - _db = db; + _tradeLog = tradeLog; _logger = logger; _jobStatus = new JobStatusRow @@ -43,10 +42,7 @@ namespace PolyTraderSharp.Services _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); + _tradeLog.EnsureIndexes(); // We do a loop waiting for items in the channel await foreach(var trade in _tradeReader.ReadAllAsync(stoppingToken)) @@ -68,15 +64,14 @@ namespace PolyTraderSharp.Services // 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) + 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 (DB-ID: {existing.TradeId})."); + _logger.Debug($"Duplikat ignoriert: ClosedTrade für Account {trade.AccountId} + Token {trade.TokenId.Substring(0, Math.Min(10, trade.TokenId.Length))}... existiert bereits."); continue; } } - - col.Insert(trade); + + _tradeLog.Insert(trade); _logger.Debug($"Saved ClosedTrade {trade.TradeId} to MongoDB"); _jobStatus.LastRun = DateTime.Now; diff --git a/services/PolymarketWssClient.cs b/services/PolymarketWssClient.cs index ff68cd0..fd8b26b 100644 --- a/services/PolymarketWssClient.cs +++ b/services/PolymarketWssClient.cs @@ -1,7 +1,6 @@ using System; -using MongoDB.Driver; using PolyTrader.Core.Persistence; -using PolyTraderSharp.Extensions; +using PolyTrader.Modules.CopyTrading.Persistence; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; @@ -24,7 +23,7 @@ namespace PolyTraderSharp.Services private readonly ServerSettings _settings; private readonly PolymarketClobClient _clob; private readonly TerminalLogger _logger; - private readonly IMongoDatabase _db; + private readonly ICopyTradeLogRepository _tradeLog; private readonly IPositionRepository _positionRepo; private readonly IAccountRepository _accountRepo; @@ -37,7 +36,7 @@ namespace PolyTraderSharp.Services ServerSettings settings, PolymarketClobClient clob, TerminalLogger logger, - IMongoDatabase db, + ICopyTradeLogRepository tradeLog, IPositionRepository positionRepo, IAccountRepository accountRepo) { @@ -46,7 +45,7 @@ namespace PolyTraderSharp.Services _settings = settings; _clob = clob; _logger = logger; - _db = db; + _tradeLog = tradeLog; _positionRepo = positionRepo; _accountRepo = accountRepo; } @@ -287,7 +286,7 @@ namespace PolyTraderSharp.Services ExitReason = "Pre Redeem" }; - _db.GetCollection("closed_trades").Insert(ct); + _tradeLog.Insert(ct); _accountRepo.Upsert(acc); _logger.Trade($"✅ [AUTO REDEEM DEMO ERFOLGREICH] {pos.MarketQuestion} | Exit: {pos.Size:F2} @ {exactLimitPrice:F3} | PnL: ${realizedPnl:F2}"); diff --git a/services/TraderAnalyticsJob.cs b/services/TraderAnalyticsJob.cs index 1576119..a4da67d 100644 --- a/services/TraderAnalyticsJob.cs +++ b/services/TraderAnalyticsJob.cs @@ -1,6 +1,5 @@ using System; -using MongoDB.Driver; -using PolyTraderSharp.Extensions; +using PolyTrader.Modules.CopyTrading.Persistence; using System.Collections.Generic; using System.Linq; using System.Threading; @@ -15,15 +14,15 @@ namespace PolyTraderSharp.Services private readonly TradingState _state; private readonly CopyTradingState _copyState; private readonly TerminalLogger _logger; - private readonly IMongoDatabase _db; + private readonly ICopyTradeLogRepository _tradeLog; private readonly JobStatusRow _jobStatus; - public TraderAnalyticsJob(TradingState state, CopyTradingState copyState, TerminalLogger logger, IMongoDatabase db, JobManager jobManager) + public TraderAnalyticsJob(TradingState state, CopyTradingState copyState, TerminalLogger logger, ICopyTradeLogRepository tradeLog, JobManager jobManager) { _state = state; _copyState = copyState; _logger = logger; - _db = db; + _tradeLog = tradeLog; _jobStatus = new JobStatusRow { @@ -86,10 +85,7 @@ namespace PolyTraderSharp.Services { _logger.Info("🔄 Starte Trader Analytics (7D / Letzte 30 Trades)..."); - var closedTradesColl = _db.GetCollection("closed_trades"); - // Ensure indexes - closedTradesColl.EnsureIndex(x => x.AccountId); - closedTradesColl.EnsureIndex(x => x.SourceTraderId); + _tradeLog.EnsureIndexes(); DateTime sevenDaysAgo = DateTime.UtcNow.AddDays(-7); @@ -102,7 +98,7 @@ namespace PolyTraderSharp.Services // The requirement says: "Welche Trades ... in den letzten 7 Tagen kopiert ... und wie hoch war die Winrate der letzten 30 Trades" // Thus we only care about MTs that had at least 1 trade in the last 7 days! int accId = acc.AccountId; - var recentMTs = closedTradesColl.LiteFind(x => x.AccountId == accId && x.ClosedAt >= sevenDaysAgo) + var recentMTs = _tradeLog.Find(x => x.AccountId == accId && x.ClosedAt >= sevenDaysAgo) .Select(x => x.SourceTraderId) .Distinct() .Where(id => id != 0) // Ignore orphaned historical trades (API resolved/auto-redeem before ID tracking patch) @@ -115,10 +111,10 @@ namespace PolyTraderSharp.Services string address = mtInfo?.WalletAddress ?? ""; // 1. Trades im 7D Fenster zählen - int trades7D = closedTradesColl.LiteFind(x => x.AccountId == accId && x.SourceTraderId == mtId && x.ClosedAt >= sevenDaysAgo).Count(); + int trades7D = _tradeLog.Find(x => x.AccountId == accId && x.SourceTraderId == mtId && x.ClosedAt >= sevenDaysAgo).Count(); // 2. Letzte 30 Trades holen - var last30 = closedTradesColl.LiteFind(x => x.AccountId == accId && x.SourceTraderId == mtId) + var last30 = _tradeLog.Find(x => x.AccountId == accId && x.SourceTraderId == mtId) .OrderByDescending(x => x.ClosedAt) .Take(30) .ToList(); diff --git a/services/TraderMonitorService.cs b/services/TraderMonitorService.cs index e2458ce..19ca230 100644 --- a/services/TraderMonitorService.cs +++ b/services/TraderMonitorService.cs @@ -1,7 +1,6 @@ using System; -using MongoDB.Driver; using PolyTrader.Core.Persistence; -using PolyTraderSharp.Extensions; +using PolyTrader.Modules.CopyTrading.Persistence; using System.Collections.Concurrent; using System.Linq; using System.Text.Json; @@ -22,7 +21,7 @@ namespace PolyTraderSharp.Services private readonly ChannelWriter _signalWriter; private readonly ChannelWriter _closedTradeWriter; private readonly TerminalLogger _logger; - private readonly IMongoDatabase? _db; + private readonly ICopyTradeLogRepository _tradeLog; private readonly IPositionRepository _positionRepo; private readonly IMarketRepository _marketRepo; @@ -49,7 +48,7 @@ namespace PolyTraderSharp.Services TerminalLogger logger, IPositionRepository positionRepo, IMarketRepository marketRepo, - IMongoDatabase? db = null) + ICopyTradeLogRepository tradeLog) { _state = state; _copyState = copyState; @@ -60,7 +59,7 @@ namespace PolyTraderSharp.Services _logger = logger; _positionRepo = positionRepo; _marketRepo = marketRepo; - _db = db; + _tradeLog = tradeLog; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) @@ -70,28 +69,24 @@ namespace PolyTraderSharp.Services // ===== STARTUP: _processedClosures aus DB vorladen ===== // Verhindert, dass nach einem Neustart alle historischen geschlossenen Trades // erneut als ClosedTrade-Records in die Datenbank geschrieben werden. - if (_db != null) + try { - try + var allClosed = _tradeLog.Find(x => !x.IsDemo); + int preloaded = 0; + foreach (var ct in allClosed) { - var closedCol = _db.GetCollection("closed_trades"); - var allClosed = closedCol.LiteFind(x => !x.IsDemo); - int preloaded = 0; - foreach (var ct in allClosed) + if (!string.IsNullOrEmpty(ct.TokenId)) { - if (!string.IsNullOrEmpty(ct.TokenId)) - { - string key = $"{ct.AccountId}_{ct.TokenId}"; - _processedClosures.TryAdd(key, true); - preloaded++; - } + string key = $"{ct.AccountId}_{ct.TokenId}"; + _processedClosures.TryAdd(key, true); + preloaded++; } - _logger.Info($"_processedClosures vorgeladen: {preloaded} Einträge aus closed_trades geladen (verhindert Duplikate nach Neustart)."); - } - catch (Exception ex) - { - _logger.Warning($"_processedClosures Preload fehlgeschlagen: {ex.Message}"); } + _logger.Info($"_processedClosures vorgeladen: {preloaded} Einträge aus closed_trades geladen (verhindert Duplikate nach Neustart)."); + } + catch (Exception ex) + { + _logger.Warning($"_processedClosures Preload fehlgeschlagen: {ex.Message}"); } @@ -337,15 +332,12 @@ namespace PolyTraderSharp.Services // Update Cache for 0ms next time _state.MarketCache[signal.TokenId] = coldItem; - // Asynchronously persist to LiteDB snapshot without blocking Hot Path - if (_db != null) - { - _ = Task.Run(() => { - try { - _marketRepo.Upsert(coldItem); - } catch { } // Failsafe - }); - } + // Asynchronously persist market snapshot without blocking Hot Path + _ = Task.Run(() => { + try { + _marketRepo.Upsert(coldItem); + } catch { } // Failsafe + }); } else { @@ -554,7 +546,7 @@ namespace PolyTraderSharp.Services // 2. Fallback: Wenn TryRemove fehlschlägt (Race Condition mit PollLiveAccountsAsync), // SourceTraderId aus der MongoDB open_positions-Tabelle wiederherstellen. - if (resolvedSourceId <= 0 && _db != null) + if (resolvedSourceId <= 0) { try { @@ -584,13 +576,8 @@ namespace PolyTraderSharp.Services } } - // 2. Prevent DB Duplicates! Fast check in MongoDB if available. - bool dbExists = false; - if (_db != null) - { - var col = _db.GetCollection("closed_trades"); - dbExists = col.Find(x => x.AccountId == acc.AccountId && x.TokenId == asset).FirstOrDefault() != null; - } + // 2. Prevent DB Duplicates! Fast check via Repository. + bool dbExists = _tradeLog.Exists(acc.AccountId, asset); string duplicateKey = $"{acc.AccountId}_{asset}"; if (!dbExists && !_processedClosures.ContainsKey(duplicateKey)) @@ -795,7 +782,7 @@ namespace PolyTraderSharp.Services DateTime resolvedOpenedAt = DateTime.UtcNow; // 4. Prüfe lokale DB-Tabelle für den Fall eines Programm-Neustarts / API-Syncs - if (resolvedTraderId == 0 && _db != null) + if (resolvedTraderId == 0) { try { @@ -860,7 +847,7 @@ namespace PolyTraderSharp.Services { // 2. Fallback: Auch im Live Sync: Wenn SourceTraderId verloren ging, // aus der MongoDB open_positions-Tabelle wiederherstellen. - if (removedPos.SourceTraderId <= 0 && _db != null) + if (removedPos.SourceTraderId <= 0) { try { diff --git a/Models/ClosedTrade.cs b/src/PolyTrader.Modules.CopyTrading/Models/ClosedTrade.cs similarity index 100% rename from Models/ClosedTrade.cs rename to src/PolyTrader.Modules.CopyTrading/Models/ClosedTrade.cs diff --git a/src/PolyTrader.Modules.CopyTrading/Persistence/ICopyTradeLogRepository.cs b/src/PolyTrader.Modules.CopyTrading/Persistence/ICopyTradeLogRepository.cs new file mode 100644 index 0000000..dfea852 --- /dev/null +++ b/src/PolyTrader.Modules.CopyTrading/Persistence/ICopyTradeLogRepository.cs @@ -0,0 +1,25 @@ +using System; +using System.Collections.Generic; +using System.Linq.Expressions; +using PolyTraderSharp.Models; + +namespace PolyTrader.Modules.CopyTrading.Persistence +{ + /// + /// Copytrading-eigener Trade-Log (Collection "closed_trades"). Enthält die kopierten, + /// geschlossenen Trades mit Copytrading-Details (SourceTraderId etc.). + /// Ergänzt später den generischen Core-Trade-Log (TradeRecord). + /// + public interface ICopyTradeLogRepository + { + void EnsureIndexes(); + + /// Prüft, ob für Account + TokenId bereits ein Eintrag existiert (Dedup). + bool Exists(int accountId, string tokenId); + + void Insert(ClosedTrade trade); + + /// Generische Abfrage (ersetzt die bisherigen LiteFind-Zugriffe). + List Find(Expression> predicate); + } +} diff --git a/src/PolyTrader.Modules.CopyTrading/Persistence/MongoCopyTradeLogRepository.cs b/src/PolyTrader.Modules.CopyTrading/Persistence/MongoCopyTradeLogRepository.cs new file mode 100644 index 0000000..39e12c6 --- /dev/null +++ b/src/PolyTrader.Modules.CopyTrading/Persistence/MongoCopyTradeLogRepository.cs @@ -0,0 +1,45 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Linq.Expressions; +using MongoDB.Driver; +using PolyTraderSharp.Models; + +namespace PolyTrader.Modules.CopyTrading.Persistence +{ + public class MongoCopyTradeLogRepository : ICopyTradeLogRepository + { + private readonly IMongoCollection _col; + + public MongoCopyTradeLogRepository(IMongoDatabase db) + { + _col = db.GetCollection("closed_trades"); + } + + public void EnsureIndexes() + { + CreateIndex(x => x.TradeId); + CreateIndex(x => x.AccountId); + CreateIndex(x => x.TokenId); + CreateIndex(x => x.SourceTraderId); + } + + private void CreateIndex(Expression> field) + { + try + { + var keys = Builders.IndexKeys.Ascending(field); + _col.Indexes.CreateOne(new CreateIndexModel(keys)); + } + catch { } + } + + public bool Exists(int accountId, string tokenId) => + _col.Find(x => x.AccountId == accountId && x.TokenId == tokenId).Any(); + + public void Insert(ClosedTrade trade) => _col.InsertOne(trade); + + public List Find(Expression> predicate) => + _col.Find(predicate).ToList(); + } +}