using CryptoExchange.Net.Logging.Extensions; using CryptoExchange.Net.Objects; using CryptoExchange.Net.Objects.Sockets; using CryptoExchange.Net.SharedApis; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace CryptoExchange.Net.Trackers.Klines { /// public class KlineTracker : IKlineTracker { private readonly IKlineSocketClient _socketClient; private readonly IKlineRestClient _restClient; private SyncStatus _status; private bool _startWithSnapshot; private ExchangeParameters? _exchangeParameters; /// /// The internal data structure /// protected readonly SortedDictionary _data = new SortedDictionary(); /// /// The pre-snapshot queue buffering updates received before the snapshot is set and which will be applied after the snapshot was set /// protected readonly List _preSnapshotQueue = new List(); /// /// Lock for accessing _data /// #if NET9_0_OR_GREATER private readonly Lock _lock = new Lock(); #else private readonly object _lock = new object(); #endif /// /// The last time the window was applied /// protected DateTime _lastWindowApplied = DateTime.MinValue; /// /// Whether or not the data has changed since last window was applied /// protected bool _changed = false; /// /// Whether the snapshot has been set /// protected bool _snapshotSet; /// /// Logger /// protected readonly ILogger _logger; /// /// Update subscription /// protected UpdateSubscription? _updateSubscription; /// /// The timestamp of the first item /// protected DateTime? _firstTimestamp; /// /// The kline interval /// public SharedKlineInterval Interval { get; } /// public SyncStatus Status { get => _status; set { if (value == _status) return; var old = _status; _status = value; _logger.KlineTrackerStatusChanged(SymbolName, old, value); OnStatusChanged?.Invoke(old, _status); } } /// public string Exchange { get; } /// public string SymbolName { get; } /// public SharedSymbol Symbol { get; } /// public int? Limit { get; } /// public TimeSpan? Period { get; } /// public DateTime? SyncedFrom { get { if (Period == null) return _firstTimestamp; var max = DateTime.UtcNow - Period.Value; if (_firstTimestamp > max) return _firstTimestamp; return max; } } /// public int Count { get { lock (_lock) { ApplyWindow(true); return _data.Count; } } } /// public SharedKline? Last { get { lock (_lock) { ApplyWindow(true); return _data.LastOrDefault().Value; } } } /// public event Func? OnAdded; /// public event Func? OnUpdated; /// public event Func? OnRemoved; /// public event Func? OnStatusChanged; /// /// ctor /// public KlineTracker( ILogger? logger, IKlineRestClient restClient, IKlineSocketClient socketClient, SharedSymbol symbol, SharedKlineInterval interval, int? limit = null, TimeSpan? period = null, ExchangeParameters? exchangeParameters = null) { _logger = logger ?? new NullLogger(); _exchangeParameters = exchangeParameters; Symbol = symbol; SymbolName = socketClient.FormatSymbol(symbol.BaseAsset, symbol.QuoteAsset, symbol.TradingMode, symbol.DeliverTime); Exchange = restClient.Exchange; Limit = limit; Period = period; Interval = interval; _socketClient = socketClient; _restClient = restClient; } /// public async Task StartAsync(bool startWithSnapshot = true) { if (Status != SyncStatus.Disconnected) throw new InvalidOperationException($"Can't start syncing unless state is {SyncStatus.Disconnected}. Current state: {Status}"); _startWithSnapshot = startWithSnapshot; Status = SyncStatus.Syncing; _logger.KlineTrackerStarting(SymbolName); var subResult = await _socketClient.SubscribeToKlineUpdatesAsync(new SubscribeKlineRequest(Symbol, Interval, exchangeParameters: _exchangeParameters), update => { AddOrUpdate(update.Data); }).ConfigureAwait(false); if (!subResult.Success) { _logger.KlineTrackerStartFailed(SymbolName, subResult.Error!.Message ?? subResult.Error!.ErrorDescription!, subResult.Error.Exception); Status = SyncStatus.Disconnected; return CallResult.Fail(subResult.Error!); } _updateSubscription = subResult.Data; _updateSubscription.ConnectionLost += HandleConnectionLost; _updateSubscription.ConnectionClosed += HandleConnectionClosed; _updateSubscription.ConnectionRestored += HandleConnectionRestored; var startResult = await DoStartAsync().ConfigureAwait(false); if (!startResult.Success) { _ = subResult.Data.CloseAsync(); Status = SyncStatus.Disconnected; return CallResult.Fail(startResult.Error!); } Status = SyncStatus.Synced; _logger.KlineTrackerStarted(SymbolName); return CallResult.Ok(); } /// public async Task StopAsync() { _logger.KlineTrackerStopping(SymbolName); Status = SyncStatus.Disconnected; await DoStopAsync().ConfigureAwait(false); _data.Clear(); _preSnapshotQueue.Clear(); _logger.KlineTrackerStopped(SymbolName); } /// /// The start procedure needed for kline syncing, generally subscribing to an update stream and requesting the snapshot /// /// protected virtual async Task DoStartAsync() { if (!_startWithSnapshot) return CallResult.Ok(); var startTime = Period == null ? (DateTime?)null : DateTime.UtcNow.Add(-Period.Value); if (_restClient.GetKlinesOptions.MaxAge != null && DateTime.UtcNow.Add(-_restClient.GetKlinesOptions.MaxAge.Value) > startTime) startTime = DateTime.UtcNow.Add(-_restClient.GetKlinesOptions.MaxAge.Value); var limit = Math.Min(_restClient.GetKlinesOptions.MaxLimit, Limit ?? 100); var request = new GetKlinesRequest(Symbol, Interval, startTime, DateTime.UtcNow, limit: limit, exchangeParameters: _exchangeParameters); var data = new List(); await foreach (var result in ExchangeHelpers.ExecutePages(_restClient.GetKlinesAsync, request).ConfigureAwait(false)) { if (!result.Success) return CallResult.Fail(result.Error!); if (Limit != null && data.Count > Limit) break; data.AddRange(result.Data); } SetInitialData(data); return CallResult.Ok(); } /// /// The stop procedure needed, generally stopping the update stream /// /// protected virtual Task DoStopAsync() => _updateSubscription?.CloseAsync() ?? Task.CompletedTask; /// public KlinesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null) { var compareTime = SyncedFrom?.AddSeconds(-2); var stats = GetStats(GetData(fromTimestamp, toTimestamp)); stats.Complete = (fromTimestamp == null || fromTimestamp >= compareTime) && (toTimestamp == null || toTimestamp >= compareTime); return stats; } private KlinesStats GetStats(IEnumerable klines) { if (!klines.Any()) return new KlinesStats(); return new KlinesStats { KlineCount = klines.Count(), FirstOpenTime = klines.First().OpenTime, LastOpenTime = klines.Last().OpenTime, HighPrice = klines.Select(d => d.LowPrice).Max(), LowPrice = klines.Select(d => d.HighPrice).Min(), Volume = klines.Select(d => d.Volume).Sum(), AverageVolume = Math.Round(klines.OrderByDescending(d => d.OpenTime).Skip(1).Select(d => d.Volume).DefaultIfEmpty().Average(), 8) }; } /// public SharedKline[] GetData(DateTime? since = null, DateTime? until = null) { lock (_lock) { ApplyWindow(true); IEnumerable result = _data.Values; if (since != null) result = result.Where(d => d.OpenTime >= since); if (until != null) result = result.Where(d => d.OpenTime <= until); return result.ToArray(); } } /// /// Set the initial kline data snapshot /// /// protected void SetInitialData(IEnumerable data) { lock (_lock) { _data.Clear(); IEnumerable items = data.OrderByDescending(d => d.OpenTime); if (Limit != null) items = items.Take(Limit.Value); if (Period != null) items = items.Where(e => e.OpenTime >= DateTime.UtcNow.Add(-Period.Value)); foreach (var item in items.OrderBy(d => d.OpenTime)) _data.Add(item.OpenTime, item); _snapshotSet = true; foreach (var item in _preSnapshotQueue) { if (_data.ContainsKey(item.OpenTime)) continue; _data.Add(item.OpenTime, item); } _firstTimestamp = _data.Count == 0 ? null : _data.Min(v => v.Key); ApplyWindow(false); _logger.KlineTrackerInitialDataSet(SymbolName, _data.Last().Key); } } /// /// Add or update a kline /// /// protected void AddOrUpdate(SharedKline item) => AddOrUpdate(new[] { item }); /// /// Add or update klines /// /// protected void AddOrUpdate(IEnumerable items) { lock (_lock) { if (_restClient != null && _startWithSnapshot && !_snapshotSet) { _preSnapshotQueue.AddRange(items); return; } foreach (var item in items) { if (_data.TryGetValue(item.OpenTime, out var existing)) { _data.Remove(item.OpenTime); _data.Add(item.OpenTime, item); OnUpdated?.Invoke(item); _logger.KlineTrackerKlineUpdated(SymbolName, _data.Last().Key); } else { _data.Add(item.OpenTime, item); OnAdded?.Invoke(item); _logger.KlineTrackerKlineAdded(SymbolName, _data.Last().Key); } } _firstTimestamp = _data.Count == 0 ? null : _data.Min(x => x.Key); _changed = true; SetSyncStatus(); ApplyWindow(true); } } private void ApplyWindow(bool broadcastEvents) { if (!_changed && (DateTime.UtcNow - _lastWindowApplied) < TimeSpan.FromSeconds(1)) return; if (Period != null) { var compareDate = DateTime.UtcNow.Add(-Period.Value); for (var i = 0; i < _data.Count; i++) { var item = _data.ElementAt(0); if (item.Key >= compareDate) break; _data.Remove(item.Key); if (broadcastEvents) OnRemoved?.Invoke(item.Value); } } if (Limit != null && _data.Count > Limit.Value) { var toRemove = Math.Max(0, _data.Count - Limit.Value); for (var i = 0; i < toRemove; i++) { var item = _data.ElementAt(0); _data.Remove(item.Key); if (broadcastEvents) OnRemoved?.Invoke(item.Value); } } _lastWindowApplied = DateTime.UtcNow; _changed = false; } private void HandleConnectionLost() { _logger.KlineTrackerConnectionLost(SymbolName); if (Status != SyncStatus.Disconnected) { Status = SyncStatus.Syncing; _snapshotSet = false; _firstTimestamp = null; _preSnapshotQueue.Clear(); } } private void HandleConnectionClosed() { _logger.KlineTrackerConnectionClosed(SymbolName); Status = SyncStatus.Disconnected; _ = StopAsync(); } private async void HandleConnectionRestored(TimeSpan _) { Status = SyncStatus.Syncing; var success = false; while (!success) { if (Status != SyncStatus.Syncing) return; var resyncResult = await DoStartAsync().ConfigureAwait(false); success = resyncResult.Success; } _logger.KlineTrackerConnectionRestored(SymbolName); SetSyncStatus(); } private void SetSyncStatus() { if (Status == SyncStatus.Synced) return; if (Period != null) { if (_firstTimestamp <= DateTime.UtcNow - Period.Value) Status = SyncStatus.Synced; else Status = SyncStatus.PartiallySynced; } if (Limit != null) { if (_data.Count == Limit.Value) Status = SyncStatus.Synced; else Status = SyncStatus.PartiallySynced; } if (Period == null && Limit == null) Status = SyncStatus.Synced; } } }