diff --git a/services/MarketSyncService.cs b/services/MarketSyncService.cs index dc20772..cde0179 100644 --- a/services/MarketSyncService.cs +++ b/services/MarketSyncService.cs @@ -3,6 +3,7 @@ using System.Threading.Tasks; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using MongoDB.Driver; +using PolyTrader.Core.Persistence; using PolyTraderSharp.Extensions; using PolyTraderSharp.Models; using PolyTraderSharp.Services; @@ -12,15 +13,15 @@ namespace PolyTraderSharp.Services public class MarketSyncService : BackgroundService { private readonly PolymarketApiService _apiService; - private readonly IMongoDatabase _db; + private readonly IMarketRepository _marketRepo; private readonly TerminalLogger _logger; private readonly JobStatusRow _jobStatus; private readonly TradingState _state; - public MarketSyncService(PolymarketApiService apiService, IMongoDatabase db, TerminalLogger logger, JobManager jobManager, TradingState state) + public MarketSyncService(PolymarketApiService apiService, IMarketRepository marketRepo, TerminalLogger logger, JobManager jobManager, TradingState state) { _apiService = apiService; - _db = db; + _marketRepo = marketRepo; _logger = logger; _state = state; @@ -93,18 +94,17 @@ namespace PolyTraderSharp.Services return; } - var col = _db.GetCollection("markets"); - col.EnsureIndex(x => x.Id); - + _marketRepo.EnsureIndexes(); + int inserted = 0; int updated = 0; - + foreach (var market in newMarkets) { - var existing = col.LiteFindOne(x => x.Id == market.Id); + var existing = _marketRepo.GetById(market.Id); if (existing == null) { - col.Insert(market); + _marketRepo.Insert(market); inserted++; // NEW: Hot-Load active markets directly into RAM Cache @@ -137,7 +137,7 @@ namespace PolyTraderSharp.Services existing.ClobTokenIds = market.ClobTokenIds; } - col.Update(existing); + _marketRepo.Update(existing); updated++; // Keep RAM cache synchronized to prevent using stale active/closed flags diff --git a/services/PolymarketWssClient.cs b/services/PolymarketWssClient.cs index 21b5db6..eafec52 100644 --- a/services/PolymarketWssClient.cs +++ b/services/PolymarketWssClient.cs @@ -1,5 +1,6 @@ using System; using MongoDB.Driver; +using PolyTrader.Core.Persistence; using PolyTraderSharp.Extensions; using System.Collections.Concurrent; using System.Collections.Generic; @@ -23,7 +24,9 @@ namespace PolyTraderSharp.Services private readonly PolymarketClobClient _clob; private readonly TerminalLogger _logger; private readonly IMongoDatabase _db; - + private readonly IPositionRepository _positionRepo; + private readonly IAccountRepository _accountRepo; + // Tracking rate limits for auto redeem: max 2 attempts per position, 5 min apart private readonly ConcurrentDictionary _redeemAttempts = new(); @@ -32,13 +35,17 @@ namespace PolyTraderSharp.Services ServerSettings settings, PolymarketClobClient clob, TerminalLogger logger, - IMongoDatabase db) + IMongoDatabase db, + IPositionRepository positionRepo, + IAccountRepository accountRepo) { _state = state; _settings = settings; _clob = clob; _logger = logger; _db = db; + _positionRepo = positionRepo; + _accountRepo = accountRepo; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) @@ -248,7 +255,7 @@ namespace PolyTraderSharp.Services { if (acc.OpenPositions.TryRemove(pos.TokenId, out _)) { - _db.GetCollection($"demo_positions_{acc.AccountId}").Delete(pos.TokenId); + _positionRepo.DeleteDemo(acc.AccountId, pos.TokenId); decimal exactLimitPrice = acc.PreRedeemLimit; decimal exitUsd = pos.Size * exactLimitPrice; @@ -278,7 +285,7 @@ namespace PolyTraderSharp.Services }; _db.GetCollection("closed_trades").Insert(ct); - _db.GetCollection("accounts").Upsert(acc); + _accountRepo.Upsert(acc); _logger.Trade($"✅ [AUTO REDEEM DEMO ERFOLGREICH] {pos.MarketQuestion} | Exit: {pos.Size:F2} @ {exactLimitPrice:F3} | PnL: ${realizedPnl:F2}"); }