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 <noreply@anthropic.com>
This commit is contained in:
Richard
2026-07-01 22:27:22 +02:00
co-authored by Claude Opus 4.8
parent 3be75c0f05
commit 88fc982089
8 changed files with 122 additions and 73 deletions
+2
View File
@@ -10,6 +10,7 @@ using Microsoft.Extensions.Options;
using PolyTrader.Core.Configuration; using PolyTrader.Core.Configuration;
using PolyTrader.Core.DependencyInjection; using PolyTrader.Core.DependencyInjection;
using PolyTrader.Core.Streaming; using PolyTrader.Core.Streaming;
using PolyTrader.Modules.CopyTrading.Persistence;
using PolyTraderSharp.Models; using PolyTraderSharp.Models;
using PolyTraderSharp.Services; using PolyTraderSharp.Services;
@@ -36,6 +37,7 @@ internal static class Program
return client.GetDatabase(dbOptions.DatabaseName); return client.GetDatabase(dbOptions.DatabaseName);
}); });
services.AddCorePersistence(); services.AddCorePersistence();
services.AddSingleton<ICopyTradeLogRepository, MongoCopyTradeLogRepository>();
services.AddSingleton<IBlockchainWssClientFactory, AlchemyWssClientFactory>(); services.AddSingleton<IBlockchainWssClientFactory, AlchemyWssClientFactory>();
services.AddSingleton((IServiceProvider sp) => ServerSettings.Load("server_settings.xml")); services.AddSingleton((IServiceProvider sp) => ServerSettings.Load("server_settings.xml"));
services.AddSingleton<TradingState>(); services.AddSingleton<TradingState>();
+8 -13
View File
@@ -1,8 +1,7 @@
using System.Threading.Channels; using System.Threading.Channels;
using MongoDB.Driver;
using PolyTraderSharp.Extensions;
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging;
using PolyTrader.Modules.CopyTrading.Persistence;
using PolyTraderSharp.Models; using PolyTraderSharp.Models;
namespace PolyTraderSharp.Services namespace PolyTraderSharp.Services
@@ -10,14 +9,14 @@ namespace PolyTraderSharp.Services
public class PersistenceService : BackgroundService public class PersistenceService : BackgroundService
{ {
private readonly ChannelReader<ClosedTrade> _tradeReader; private readonly ChannelReader<ClosedTrade> _tradeReader;
private readonly IMongoDatabase _db; private readonly ICopyTradeLogRepository _tradeLog;
private readonly TerminalLogger _logger; private readonly TerminalLogger _logger;
private readonly JobStatusRow _jobStatus; private readonly JobStatusRow _jobStatus;
public PersistenceService(ChannelReader<ClosedTrade> tradeReader, IMongoDatabase db, TerminalLogger logger, JobManager jobManager) public PersistenceService(ChannelReader<ClosedTrade> tradeReader, ICopyTradeLogRepository tradeLog, TerminalLogger logger, JobManager jobManager)
{ {
_tradeReader = tradeReader; _tradeReader = tradeReader;
_db = db; _tradeLog = tradeLog;
_logger = logger; _logger = logger;
_jobStatus = new JobStatusRow _jobStatus = new JobStatusRow
@@ -43,10 +42,7 @@ namespace PolyTraderSharp.Services
_jobStatus.StatusText = "Listening (Channel)..."; _jobStatus.StatusText = "Listening (Channel)...";
// One-time index setup (moved out of hot loop) // One-time index setup (moved out of hot loop)
var col = _db.GetCollection<ClosedTrade>("closed_trades"); _tradeLog.EnsureIndexes();
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 // We do a loop waiting for items in the channel
await foreach(var trade in _tradeReader.ReadAllAsync(stoppingToken)) 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. // Trades bei jedem Sync-Zyklus oder nach einem Neustart erneut eingefügt werden.
if (!string.IsNullOrEmpty(trade.TokenId)) if (!string.IsNullOrEmpty(trade.TokenId))
{ {
var existing = col.LiteFindOne(x => x.AccountId == trade.AccountId && x.TokenId == trade.TokenId); if (_tradeLog.Exists(trade.AccountId, 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})."); _logger.Debug($"Duplikat ignoriert: ClosedTrade für Account {trade.AccountId} + Token {trade.TokenId.Substring(0, Math.Min(10, trade.TokenId.Length))}... existiert bereits.");
continue; continue;
} }
} }
col.Insert(trade); _tradeLog.Insert(trade);
_logger.Debug($"Saved ClosedTrade {trade.TradeId} to MongoDB"); _logger.Debug($"Saved ClosedTrade {trade.TradeId} to MongoDB");
_jobStatus.LastRun = DateTime.Now; _jobStatus.LastRun = DateTime.Now;
+5 -6
View File
@@ -1,7 +1,6 @@
using System; using System;
using MongoDB.Driver;
using PolyTrader.Core.Persistence; using PolyTrader.Core.Persistence;
using PolyTraderSharp.Extensions; using PolyTrader.Modules.CopyTrading.Persistence;
using System.Collections.Concurrent; using System.Collections.Concurrent;
using System.Collections.Generic; using System.Collections.Generic;
using System.Linq; using System.Linq;
@@ -24,7 +23,7 @@ namespace PolyTraderSharp.Services
private readonly ServerSettings _settings; private readonly ServerSettings _settings;
private readonly PolymarketClobClient _clob; private readonly PolymarketClobClient _clob;
private readonly TerminalLogger _logger; private readonly TerminalLogger _logger;
private readonly IMongoDatabase _db; private readonly ICopyTradeLogRepository _tradeLog;
private readonly IPositionRepository _positionRepo; private readonly IPositionRepository _positionRepo;
private readonly IAccountRepository _accountRepo; private readonly IAccountRepository _accountRepo;
@@ -37,7 +36,7 @@ namespace PolyTraderSharp.Services
ServerSettings settings, ServerSettings settings,
PolymarketClobClient clob, PolymarketClobClient clob,
TerminalLogger logger, TerminalLogger logger,
IMongoDatabase db, ICopyTradeLogRepository tradeLog,
IPositionRepository positionRepo, IPositionRepository positionRepo,
IAccountRepository accountRepo) IAccountRepository accountRepo)
{ {
@@ -46,7 +45,7 @@ namespace PolyTraderSharp.Services
_settings = settings; _settings = settings;
_clob = clob; _clob = clob;
_logger = logger; _logger = logger;
_db = db; _tradeLog = tradeLog;
_positionRepo = positionRepo; _positionRepo = positionRepo;
_accountRepo = accountRepo; _accountRepo = accountRepo;
} }
@@ -287,7 +286,7 @@ namespace PolyTraderSharp.Services
ExitReason = "Pre Redeem" ExitReason = "Pre Redeem"
}; };
_db.GetCollection<ClosedTrade>("closed_trades").Insert(ct); _tradeLog.Insert(ct);
_accountRepo.Upsert(acc); _accountRepo.Upsert(acc);
_logger.Trade($"✅ [AUTO REDEEM DEMO ERFOLGREICH] {pos.MarketQuestion} | Exit: {pos.Size:F2} @ {exactLimitPrice:F3} | PnL: ${realizedPnl:F2}"); _logger.Trade($"✅ [AUTO REDEEM DEMO ERFOLGREICH] {pos.MarketQuestion} | Exit: {pos.Size:F2} @ {exactLimitPrice:F3} | PnL: ${realizedPnl:F2}");
+8 -12
View File
@@ -1,6 +1,5 @@
using System; using System;
using MongoDB.Driver; using PolyTrader.Modules.CopyTrading.Persistence;
using PolyTraderSharp.Extensions;
using System.Collections.Generic; using System.Collections.Generic;
using System.Linq; using System.Linq;
using System.Threading; using System.Threading;
@@ -15,15 +14,15 @@ namespace PolyTraderSharp.Services
private readonly TradingState _state; private readonly TradingState _state;
private readonly CopyTradingState _copyState; private readonly CopyTradingState _copyState;
private readonly TerminalLogger _logger; private readonly TerminalLogger _logger;
private readonly IMongoDatabase _db; private readonly ICopyTradeLogRepository _tradeLog;
private readonly JobStatusRow _jobStatus; 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; _state = state;
_copyState = copyState; _copyState = copyState;
_logger = logger; _logger = logger;
_db = db; _tradeLog = tradeLog;
_jobStatus = new JobStatusRow _jobStatus = new JobStatusRow
{ {
@@ -86,10 +85,7 @@ namespace PolyTraderSharp.Services
{ {
_logger.Info("🔄 Starte Trader Analytics (7D / Letzte 30 Trades)..."); _logger.Info("🔄 Starte Trader Analytics (7D / Letzte 30 Trades)...");
var closedTradesColl = _db.GetCollection<ClosedTrade>("closed_trades"); _tradeLog.EnsureIndexes();
// Ensure indexes
closedTradesColl.EnsureIndex(x => x.AccountId);
closedTradesColl.EnsureIndex(x => x.SourceTraderId);
DateTime sevenDaysAgo = DateTime.UtcNow.AddDays(-7); 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" // 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! // Thus we only care about MTs that had at least 1 trade in the last 7 days!
int accId = acc.AccountId; 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) .Select(x => x.SourceTraderId)
.Distinct() .Distinct()
.Where(id => id != 0) // Ignore orphaned historical trades (API resolved/auto-redeem before ID tracking patch) .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 ?? ""; string address = mtInfo?.WalletAddress ?? "";
// 1. Trades im 7D Fenster zählen // 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 // 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) .OrderByDescending(x => x.ClosedAt)
.Take(30) .Take(30)
.ToList(); .ToList();
+11 -24
View File
@@ -1,7 +1,6 @@
using System; using System;
using MongoDB.Driver;
using PolyTrader.Core.Persistence; using PolyTrader.Core.Persistence;
using PolyTraderSharp.Extensions; using PolyTrader.Modules.CopyTrading.Persistence;
using System.Collections.Concurrent; using System.Collections.Concurrent;
using System.Linq; using System.Linq;
using System.Text.Json; using System.Text.Json;
@@ -22,7 +21,7 @@ namespace PolyTraderSharp.Services
private readonly ChannelWriter<CopySignal> _signalWriter; private readonly ChannelWriter<CopySignal> _signalWriter;
private readonly ChannelWriter<ClosedTrade> _closedTradeWriter; private readonly ChannelWriter<ClosedTrade> _closedTradeWriter;
private readonly TerminalLogger _logger; private readonly TerminalLogger _logger;
private readonly IMongoDatabase? _db; private readonly ICopyTradeLogRepository _tradeLog;
private readonly IPositionRepository _positionRepo; private readonly IPositionRepository _positionRepo;
private readonly IMarketRepository _marketRepo; private readonly IMarketRepository _marketRepo;
@@ -49,7 +48,7 @@ namespace PolyTraderSharp.Services
TerminalLogger logger, TerminalLogger logger,
IPositionRepository positionRepo, IPositionRepository positionRepo,
IMarketRepository marketRepo, IMarketRepository marketRepo,
IMongoDatabase? db = null) ICopyTradeLogRepository tradeLog)
{ {
_state = state; _state = state;
_copyState = copyState; _copyState = copyState;
@@ -60,7 +59,7 @@ namespace PolyTraderSharp.Services
_logger = logger; _logger = logger;
_positionRepo = positionRepo; _positionRepo = positionRepo;
_marketRepo = marketRepo; _marketRepo = marketRepo;
_db = db; _tradeLog = tradeLog;
} }
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
@@ -70,12 +69,9 @@ namespace PolyTraderSharp.Services
// ===== STARTUP: _processedClosures aus DB vorladen ===== // ===== STARTUP: _processedClosures aus DB vorladen =====
// Verhindert, dass nach einem Neustart alle historischen geschlossenen Trades // Verhindert, dass nach einem Neustart alle historischen geschlossenen Trades
// erneut als ClosedTrade-Records in die Datenbank geschrieben werden. // erneut als ClosedTrade-Records in die Datenbank geschrieben werden.
if (_db != null)
{
try try
{ {
var closedCol = _db.GetCollection<ClosedTrade>("closed_trades"); var allClosed = _tradeLog.Find(x => !x.IsDemo);
var allClosed = closedCol.LiteFind(x => !x.IsDemo);
int preloaded = 0; int preloaded = 0;
foreach (var ct in allClosed) foreach (var ct in allClosed)
{ {
@@ -92,7 +88,6 @@ namespace PolyTraderSharp.Services
{ {
_logger.Warning($"_processedClosures Preload fehlgeschlagen: {ex.Message}"); _logger.Warning($"_processedClosures Preload fehlgeschlagen: {ex.Message}");
} }
}
// ===== STARTUP: Sofortiger Warmup der Live-Positionen und Master-Tracker ===== // ===== STARTUP: Sofortiger Warmup der Live-Positionen und Master-Tracker =====
@@ -337,16 +332,13 @@ namespace PolyTraderSharp.Services
// Update Cache for 0ms next time // Update Cache for 0ms next time
_state.MarketCache[signal.TokenId] = coldItem; _state.MarketCache[signal.TokenId] = coldItem;
// Asynchronously persist to LiteDB snapshot without blocking Hot Path // Asynchronously persist market snapshot without blocking Hot Path
if (_db != null)
{
_ = Task.Run(() => { _ = Task.Run(() => {
try { try {
_marketRepo.Upsert(coldItem); _marketRepo.Upsert(coldItem);
} catch { } // Failsafe } catch { } // Failsafe
}); });
} }
}
else else
{ {
// Provide fallback display if un-cached and API fails // Provide fallback display if un-cached and API fails
@@ -554,7 +546,7 @@ namespace PolyTraderSharp.Services
// 2. Fallback: Wenn TryRemove fehlschlägt (Race Condition mit PollLiveAccountsAsync), // 2. Fallback: Wenn TryRemove fehlschlägt (Race Condition mit PollLiveAccountsAsync),
// SourceTraderId aus der MongoDB open_positions-Tabelle wiederherstellen. // SourceTraderId aus der MongoDB open_positions-Tabelle wiederherstellen.
if (resolvedSourceId <= 0 && _db != null) if (resolvedSourceId <= 0)
{ {
try try
{ {
@@ -584,13 +576,8 @@ namespace PolyTraderSharp.Services
} }
} }
// 2. Prevent DB Duplicates! Fast check in MongoDB if available. // 2. Prevent DB Duplicates! Fast check via Repository.
bool dbExists = false; bool dbExists = _tradeLog.Exists(acc.AccountId, asset);
if (_db != null)
{
var col = _db.GetCollection<ClosedTrade>("closed_trades");
dbExists = col.Find(x => x.AccountId == acc.AccountId && x.TokenId == asset).FirstOrDefault() != null;
}
string duplicateKey = $"{acc.AccountId}_{asset}"; string duplicateKey = $"{acc.AccountId}_{asset}";
if (!dbExists && !_processedClosures.ContainsKey(duplicateKey)) if (!dbExists && !_processedClosures.ContainsKey(duplicateKey))
@@ -795,7 +782,7 @@ namespace PolyTraderSharp.Services
DateTime resolvedOpenedAt = DateTime.UtcNow; DateTime resolvedOpenedAt = DateTime.UtcNow;
// 4. Prüfe lokale DB-Tabelle für den Fall eines Programm-Neustarts / API-Syncs // 4. Prüfe lokale DB-Tabelle für den Fall eines Programm-Neustarts / API-Syncs
if (resolvedTraderId == 0 && _db != null) if (resolvedTraderId == 0)
{ {
try try
{ {
@@ -860,7 +847,7 @@ namespace PolyTraderSharp.Services
{ {
// 2. Fallback: Auch im Live Sync: Wenn SourceTraderId verloren ging, // 2. Fallback: Auch im Live Sync: Wenn SourceTraderId verloren ging,
// aus der MongoDB open_positions-Tabelle wiederherstellen. // aus der MongoDB open_positions-Tabelle wiederherstellen.
if (removedPos.SourceTraderId <= 0 && _db != null) if (removedPos.SourceTraderId <= 0)
{ {
try try
{ {
@@ -0,0 +1,25 @@
using System;
using System.Collections.Generic;
using System.Linq.Expressions;
using PolyTraderSharp.Models;
namespace PolyTrader.Modules.CopyTrading.Persistence
{
/// <summary>
/// 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).
/// </summary>
public interface ICopyTradeLogRepository
{
void EnsureIndexes();
/// <summary>Prüft, ob für Account + TokenId bereits ein Eintrag existiert (Dedup).</summary>
bool Exists(int accountId, string tokenId);
void Insert(ClosedTrade trade);
/// <summary>Generische Abfrage (ersetzt die bisherigen LiteFind-Zugriffe).</summary>
List<ClosedTrade> Find(Expression<Func<ClosedTrade, bool>> predicate);
}
}
@@ -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<ClosedTrade> _col;
public MongoCopyTradeLogRepository(IMongoDatabase db)
{
_col = db.GetCollection<ClosedTrade>("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<Func<ClosedTrade, object>> field)
{
try
{
var keys = Builders<ClosedTrade>.IndexKeys.Ascending(field);
_col.Indexes.CreateOne(new CreateIndexModel<ClosedTrade>(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<ClosedTrade> Find(Expression<Func<ClosedTrade, bool>> predicate) =>
_col.Find(predicate).ToList();
}
}