Phase 3d (1/2): TraderMonitorService auf Repositories umgestellt

- 8 Positions-Zugriffe (open_positions_{id}) -> IPositionRepository
  (FindLive/UpsertLive/DeleteLive).
- 1 Market-Upsert -> IMarketRepository.
- closed_trades bleibt auf _db (ClosedTrade-Repo folgt in Phase 5).
- Verhalten unverändert (Repos bilden Shim-Semantik 1:1 nach); Build 0 Fehler.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
bergm
2026-07-01 16:39:28 +02:00
co-authored by Claude Opus 4.8
parent 41c43fa3aa
commit a0e367e645
+16 -13
View File
@@ -1,5 +1,6 @@
using System; using System;
using MongoDB.Driver; using MongoDB.Driver;
using PolyTrader.Core.Persistence;
using PolyTraderSharp.Extensions; using PolyTraderSharp.Extensions;
using System.Collections.Concurrent; using System.Collections.Concurrent;
using System.Linq; using System.Linq;
@@ -21,6 +22,8 @@ namespace PolyTraderSharp.Services
private readonly ChannelWriter<ClosedTrade> _closedTradeWriter; private readonly ChannelWriter<ClosedTrade> _closedTradeWriter;
private readonly TerminalLogger _logger; private readonly TerminalLogger _logger;
private readonly IMongoDatabase? _db; private readonly IMongoDatabase? _db;
private readonly IPositionRepository _positionRepo;
private readonly IMarketRepository _marketRepo;
// Prevents duplicates. Fast O(1) lookup cache to prevent DB spam. // Prevents duplicates. Fast O(1) lookup cache to prevent DB spam.
private readonly ConcurrentDictionary<string, DateTime> _processedTxHashes = new(); private readonly ConcurrentDictionary<string, DateTime> _processedTxHashes = new();
@@ -42,6 +45,8 @@ namespace PolyTraderSharp.Services
ChannelWriter<CopySignal> signalWriter, ChannelWriter<CopySignal> signalWriter,
ChannelWriter<ClosedTrade> closedTradeWriter, ChannelWriter<ClosedTrade> closedTradeWriter,
TerminalLogger logger, TerminalLogger logger,
IPositionRepository positionRepo,
IMarketRepository marketRepo,
IMongoDatabase? db = null) IMongoDatabase? db = null)
{ {
_state = state; _state = state;
@@ -50,6 +55,8 @@ namespace PolyTraderSharp.Services
_signalWriter = signalWriter; _signalWriter = signalWriter;
_closedTradeWriter = closedTradeWriter; _closedTradeWriter = closedTradeWriter;
_logger = logger; _logger = logger;
_positionRepo = positionRepo;
_marketRepo = marketRepo;
_db = db; _db = db;
} }
@@ -332,8 +339,7 @@ namespace PolyTraderSharp.Services
{ {
_ = Task.Run(() => { _ = Task.Run(() => {
try { try {
var mdColl = _db.GetCollection<MarketData>("markets"); _marketRepo.Upsert(coldItem);
mdColl.Upsert(coldItem);
} catch { } // Failsafe } catch { } // Failsafe
}); });
} }
@@ -549,8 +555,7 @@ namespace PolyTraderSharp.Services
{ {
try try
{ {
var liveCol = _db.GetCollection<Position>($"open_positions_{acc.AccountId}"); var dbPos = _positionRepo.FindLive(acc.AccountId, asset);
var dbPos = liveCol.LiteFindOne(x => x.TokenId == asset);
if (dbPos != null && dbPos.SourceTraderId > 0) if (dbPos != null && dbPos.SourceTraderId > 0)
{ {
resolvedSourceId = dbPos.SourceTraderId; resolvedSourceId = dbPos.SourceTraderId;
@@ -705,7 +710,7 @@ namespace PolyTraderSharp.Services
existing.CurrentValueUsd = curValue; existing.CurrentValueUsd = curValue;
if (expiry.HasValue) existing.ExpiryDate = expiry; if (expiry.HasValue) existing.ExpiryDate = expiry;
if (_db != null) { try { _db.GetCollection<Position>($"open_positions_{acc.AccountId}").Upsert(existing); } catch { } } try { _positionRepo.UpsertLive(acc.AccountId, existing); } catch { }
// Auto-Redeem Fallback via REST // Auto-Redeem Fallback via REST
if (acc.PreRedeemLimit > 0 && curPrice >= acc.PreRedeemLimit && acc.IsActive) if (acc.PreRedeemLimit > 0 && curPrice >= acc.PreRedeemLimit && acc.IsActive)
@@ -791,8 +796,7 @@ namespace PolyTraderSharp.Services
{ {
try try
{ {
var liveCol = _db.GetCollection<Position>($"open_positions_{acc.AccountId}"); var dbPos = _positionRepo.FindLive(acc.AccountId, asset);
var dbPos = liveCol.LiteFindOne(x => x.TokenId == asset);
if (dbPos != null) if (dbPos != null)
{ {
if (dbPos.SourceTraderId > 0) if (dbPos.SourceTraderId > 0)
@@ -829,7 +833,7 @@ namespace PolyTraderSharp.Services
OpenedAt = resolvedOpenedAt OpenedAt = resolvedOpenedAt
}; };
acc.OpenPositions.TryAdd(asset, newPos); acc.OpenPositions.TryAdd(asset, newPos);
if (_db != null) { try { _db.GetCollection<Position>($"open_positions_{acc.AccountId}").Upsert(newPos); } catch { } } try { _positionRepo.UpsertLive(acc.AccountId, newPos); } catch { }
if (resolvedTraderId > 0) if (resolvedTraderId > 0)
_logger.Info($"🌐 Live Position erkannt: {title} ({newPos.Outcome}) - ${amountUsd} - Account: {acc.Name} [Zugeordnet: {resolvedTraderName}]"); _logger.Info($"🌐 Live Position erkannt: {title} ({newPos.Outcome}) - ${amountUsd} - Account: {acc.Name} [Zugeordnet: {resolvedTraderName}]");
@@ -857,8 +861,7 @@ namespace PolyTraderSharp.Services
{ {
try try
{ {
var liveCol = _db.GetCollection<Position>($"open_positions_{acc.AccountId}"); var dbPos = _positionRepo.FindLive(acc.AccountId, k);
var dbPos = liveCol.LiteFindOne(x => x.TokenId == k);
if (dbPos != null && dbPos.SourceTraderId > 0) if (dbPos != null && dbPos.SourceTraderId > 0)
{ {
removedPos.SourceTraderId = dbPos.SourceTraderId; removedPos.SourceTraderId = dbPos.SourceTraderId;
@@ -880,7 +883,7 @@ namespace PolyTraderSharp.Services
if (matchedClose.HasValue) if (matchedClose.HasValue)
{ {
acc.OpenPositions.TryRemove(k, out _); // Safe removal! acc.OpenPositions.TryRemove(k, out _); // Safe removal!
if (_db != null) { try { _db.GetCollection<Position>($"open_positions_{acc.AccountId}").Delete(k); } catch { } } try { _positionRepo.DeleteLive(acc.AccountId, k); } catch { }
decimal realizedPnl = 0m; decimal realizedPnl = 0m;
DateTime resolvedClosedAt = DateTime.UtcNow; DateTime resolvedClosedAt = DateTime.UtcNow;
@@ -953,7 +956,7 @@ namespace PolyTraderSharp.Services
if (isClosed) if (isClosed)
{ {
acc.OpenPositions.TryRemove(k, out _); // Safe removal! acc.OpenPositions.TryRemove(k, out _); // Safe removal!
if (_db != null) { try { _db.GetCollection<Position>($"open_positions_{acc.AccountId}").Delete(k); } catch { } } try { _positionRepo.DeleteLive(acc.AccountId, k); } catch { }
decimal exitPrice = isWinner ? 1.0m : 0.0m; decimal exitPrice = isWinner ? 1.0m : 0.0m;
decimal exitUsd = removedPos.Size * exitPrice; decimal exitUsd = removedPos.Size * exitPrice;
@@ -1029,7 +1032,7 @@ namespace PolyTraderSharp.Services
{ {
if (acc.OpenPositions.TryRemove(k, out _)) if (acc.OpenPositions.TryRemove(k, out _))
{ {
if (_db != null) { try { _db.GetCollection<Position>($"open_positions_{acc.AccountId}").Delete(k); } catch { } } try { _positionRepo.DeleteLive(acc.AccountId, k); } catch { }
_logger.Info($"🌐 Live Position {removedPos.MarketQuestion} final entfernt (Ext. Verkauft/Wartend nach {ageMinutes:F0} Min.)"); _logger.Info($"🌐 Live Position {removedPos.MarketQuestion} final entfernt (Ext. Verkauft/Wartend nach {ageMinutes:F0} Min.)");
} }
} }