1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-21 05:13:21 +00:00
Fix for intermittently failing rate limiting test
Added ConnectionId to RequestDefinition to correctly handle connection and path rate limiting configuration
Added ValidateMessage method to websocket Query object to filter messages even though it is matched to the query based on the  ListenIdentifier
Added KlineTracker and TradeTracker implementation
This commit is contained in:
Jan Korf
2024-10-28 10:36:19 +01:00
committed by GitHub
parent ed007b5272
commit 9e86a08327
17 changed files with 2386 additions and 5 deletions
@@ -0,0 +1,34 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers
{
/// <summary>
/// Compare value
/// </summary>
public record CompareValue
{
/// <summary>
/// The value difference
/// </summary>
public decimal? Difference { get; set; }
/// <summary>
/// The value difference percentage
/// </summary>
public decimal? PercentageDifference { get; set; }
/// <summary>
/// ctor
/// </summary>
public CompareValue(decimal? value1, decimal? value2)
{
if (value1 == null || value2 == null)
return;
Difference = value2 - value1;
PercentageDifference = value1.Value == 0 ? null : Math.Round(value2.Value / value1.Value * 100 - 100, 4);
}
}
}
@@ -0,0 +1,105 @@
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.SharedApis;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// A tracker for kline data of a symbol
/// </summary>
public interface IKlineTracker
{
/// <summary>
/// The total number of klines
/// </summary>
int Count { get; }
/// <summary>
/// Exchange name
/// </summary>
string Exchange { get; }
/// <summary>
/// Symbol name
/// </summary>
string SymbolName { get; }
/// <summary>
/// Symbol
/// </summary>
SharedSymbol Symbol { get; }
/// <summary>
/// The max number of klines tracked
/// </summary>
int? Limit { get; }
/// <summary>
/// The max age of the data tracked
/// </summary>
TimeSpan? Period { get; }
/// <summary>
/// From which timestamp the trades are registered
/// </summary>
DateTime? SyncedFrom { get; }
/// <summary>
/// Sync status
/// </summary>
SyncStatus Status { get; }
/// <summary>
/// Get the last kline
/// </summary>
SharedKline? Last { get; }
/// <summary>
/// Event for when a new kline is added
/// </summary>
event Func<SharedKline, Task>? OnAdded;
/// <summary>
/// Event for when a kline is removed because it's no longer within the period/limit window
/// </summary>
event Func<SharedKline, Task>? OnRemoved;
/// <summary>
/// Event for when a kline is updated
/// </summary>
event Func<SharedKline, Task> OnUpdated;
/// <summary>
/// Event for when the sync status changes
/// </summary>
event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// Start synchronization
/// </summary>
/// <returns></returns>
Task<CallResult> StartAsync(bool startWithSnapshot = true);
/// <summary>
/// Stop synchronization
/// </summary>
/// <returns></returns>
Task StopAsync();
/// <summary>
/// Get the data tracked
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
IEnumerable<SharedKline> GetData(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
/// <summary>
/// Get statitistics on the klines
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
KlinesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
}
}
@@ -0,0 +1,481 @@
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.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <inheritdoc />
public class KlineTracker : IKlineTracker
{
private readonly IKlineSocketClient _socketClient;
private readonly IKlineRestClient _restClient;
private SyncStatus _status;
private bool _startWithSnapshot;
/// <summary>
/// The internal data structure
/// </summary>
protected readonly Dictionary<DateTime, SharedKline> _data = new Dictionary<DateTime, SharedKline>();
/// <summary>
/// The pre-snapshot queue buffering updates received before the snapshot is set and which will be applied after the snapshot was set
/// </summary>
protected readonly List<SharedKline> _preSnapshotQueue = new List<SharedKline>();
/// <summary>
/// Lock for accessing _data
/// </summary>
protected readonly object _lock = new object();
/// <summary>
/// The last time the window was applied
/// </summary>
protected DateTime _lastWindowApplied = DateTime.MinValue;
/// <summary>
/// Whether or not the data has changed since last window was applied
/// </summary>
protected bool _changed = false;
/// <summary>
/// The kline interval
/// </summary>
protected readonly SharedKlineInterval _interval;
/// <summary>
/// Whether the snapshot has been set
/// </summary>
protected bool _snapshotSet;
/// <summary>
/// Logger
/// </summary>
protected readonly ILogger _logger;
/// <summary>
/// Update subscription
/// </summary>
protected UpdateSubscription? _updateSubscription;
/// <summary>
/// The timestamp of the first item
/// </summary>
protected DateTime? _firstTimestamp;
/// <inheritdoc/>
public SyncStatus Status
{
get => _status;
set
{
if (value == _status)
return;
var old = _status;
_status = value;
_logger.KlineTrackerStatusChanged(SymbolName, old, value);
OnStatusChanged?.Invoke(old, _status);
}
}
/// <inheritdoc />
public string Exchange { get; }
/// <inheritdoc />
public string SymbolName { get; }
/// <inheritdoc />
public SharedSymbol Symbol { get; }
/// <inheritdoc/>
public int? Limit { get; }
/// <inheritdoc/>
public TimeSpan? Period { get; }
/// <inheritdoc />
public DateTime? SyncedFrom
{
get
{
if (Period == null)
return _firstTimestamp;
var max = DateTime.UtcNow - Period.Value;
if (_firstTimestamp > max)
return _firstTimestamp;
return max;
}
}
/// <inheritdoc />
public int Count
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.Count;
}
}
}
/// <inheritdoc />
public SharedKline? Last
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.LastOrDefault().Value;
}
}
}
/// <inheritdoc />
public event Func<SharedKline, Task>? OnAdded;
/// <inheritdoc />
public event Func<SharedKline, Task>? OnUpdated;
/// <inheritdoc />
public event Func<SharedKline, Task>? OnRemoved;
/// <inheritdoc />
public event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// ctor
/// </summary>
public KlineTracker(
ILogger? logger,
IKlineRestClient restClient,
IKlineSocketClient socketClient,
SharedSymbol symbol,
SharedKlineInterval interval,
int? limit = null,
TimeSpan? period = null)
{
_logger = logger ?? new NullLogger<KlineTracker>();
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;
}
/// <inheritdoc />
public async Task<CallResult> 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 startResult = await DoStartAsync().ConfigureAwait(false);
if (!startResult)
{
_logger.KlineTrackerStartFailed(SymbolName, startResult.Error!.ToString());
Status = SyncStatus.Disconnected;
return new CallResult(startResult.Error!);
}
_updateSubscription = startResult.Data;
_updateSubscription.ConnectionLost += HandleConnectionLost;
_updateSubscription.ConnectionClosed += HandleConnectionClosed;
_updateSubscription.ConnectionRestored += HandleConnectionRestored;
Status = SyncStatus.Synced;
_logger.KlineTrackerStarted(SymbolName);
return new CallResult(null);
}
/// <inheritdoc />
public async Task StopAsync()
{
_logger.KlineTrackerStopping(SymbolName);
Status = SyncStatus.Disconnected;
await DoStopAsync().ConfigureAwait(false);
_data.Clear();
_preSnapshotQueue.Clear();
_logger.KlineTrackerStopped(SymbolName);
}
/// <summary>
/// The start procedure needed for kline syncing, generally subscribing to an update stream and requesting the snapshot
/// </summary>
/// <returns></returns>
protected virtual async Task<CallResult<UpdateSubscription>> DoStartAsync()
{
var subResult = await _socketClient.SubscribeToKlineUpdatesAsync(new SubscribeKlineRequest(Symbol, _interval),
update =>
{
AddOrUpdate(update.Data);
}).ConfigureAwait(false);
if (!subResult)
{
Status = SyncStatus.Disconnected;
return subResult;
}
if (!_startWithSnapshot)
return subResult;
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.MaxRequestDataPoints ?? _restClient.GetKlinesOptions.MaxTotalDataPoints ?? 100, Limit ?? 100);
var request = new GetKlinesRequest(Symbol, _interval, startTime, DateTime.UtcNow, limit: limit);
var data = new List<SharedKline>();
await foreach (var result in ExchangeHelpers.ExecutePages(_restClient.GetKlinesAsync, request).ConfigureAwait(false))
{
if (!result)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(result.Error!);
}
if (Limit != null && data.Count > Limit)
break;
data.AddRange(result.Data);
}
SetInitialData(data);
return subResult;
}
/// <summary>
/// The stop procedure needed, generally stopping the update stream
/// </summary>
/// <returns></returns>
protected virtual Task DoStopAsync() => _updateSubscription?.CloseAsync() ?? Task.CompletedTask;
/// <inheritdoc />
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<SharedKline> 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)
};
}
/// <inheritdoc />
public IEnumerable<SharedKline> GetData(DateTime? since = null, DateTime? until = null)
{
lock (_lock)
{
ApplyWindow(true);
IEnumerable<SharedKline> 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.ToList();
}
}
/// <summary>
/// Set the initial kline data snapshot
/// </summary>
/// <param name="data"></param>
protected void SetInitialData(IEnumerable<SharedKline> data)
{
lock (_lock)
{
_data.Clear();
IEnumerable<SharedKline> 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.Min(v => v.Key);
ApplyWindow(false);
_logger.KlineTrackerInitialDataSet(SymbolName, _data.Last().Key);
}
}
/// <summary>
/// Add or update a kline
/// </summary>
/// <param name="item"></param>
protected void AddOrUpdate(SharedKline item) => AddOrUpdate(new[] { item });
/// <summary>
/// Add or update klines
/// </summary>
/// <param name="items"></param>
protected void AddOrUpdate(IEnumerable<SharedKline> 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.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;
}
_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;
}
}
}
@@ -0,0 +1,30 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// Klines statistics comparison
/// </summary>
public record KlinesCompare
{
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? LowPriceDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? HighPriceDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? VolumeDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? AverageVolumeDif { get; set; }
}
}
@@ -0,0 +1,59 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// Klines statistics
/// </summary>
public record KlinesStats
{
/// <summary>
/// Number of klines
/// </summary>
public int KlineCount { get; set; }
/// <summary>
/// The kline open time of the first entry
/// </summary>
public DateTime? FirstOpenTime { get; set; }
/// <summary>
/// The kline open time of the last entry
/// </summary>
public DateTime? LastOpenTime { get; set; }
/// <summary>
/// Lowest trade price
/// </summary>
public decimal? LowPrice { get; set; }
/// <summary>
/// Highest trade price
/// </summary>
public decimal? HighPrice { get; set; }
/// <summary>
/// Trade volume
/// </summary>
public decimal Volume { get; set; }
/// <summary>
/// Average volume per kline
/// </summary>
public decimal? AverageVolume { get; set; }
/// <summary>
/// Whether the data is complete
/// </summary>
public bool Complete { get; set; }
/// <summary>
/// Compare 2 stat snapshots to eachother
/// </summary>
public KlinesCompare CompareTo(KlinesStats otherStats)
{
return new KlinesCompare
{
LowPriceDif = new CompareValue(LowPrice, otherStats.LowPrice),
HighPriceDif = new CompareValue(HighPrice, otherStats.HighPrice),
VolumeDif = new CompareValue(Volume, otherStats.Volume),
AverageVolumeDif = new CompareValue(AverageVolume, otherStats.AverageVolume),
};
}
}
}
@@ -0,0 +1,100 @@
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.SharedApis;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// A tracker for trades on a symbol
/// </summary>
public interface ITradeTracker
{
/// <summary>
/// The total number of trades
/// </summary>
int Count { get; }
/// <summary>
/// Exchange name
/// </summary>
string Exchange { get; }
/// <summary>
/// Symbol name
/// </summary>
string SymbolName { get; }
/// <summary>
/// Symbol
/// </summary>
SharedSymbol Symbol { get; }
/// <summary>
/// The max number of trades tracked
/// </summary>
int? Limit { get; }
/// <summary>
/// The max age of the data tracked
/// </summary>
TimeSpan? Period { get; }
/// <summary>
/// From which timestamp the trades are registered
/// </summary>
DateTime? SyncedFrom { get; }
/// <summary>
/// The current synchronization status
/// </summary>
SyncStatus Status { get; }
/// <summary>
/// Get the last trade
/// </summary>
SharedTrade? Last { get; }
/// <summary>
/// Event for when a new trade is added
/// </summary>
event Func<SharedTrade, Task>? OnAdded;
/// <summary>
/// Event for when a trade is removed because it's no longer within the period/limit window
/// </summary>
event Func<SharedTrade, Task>? OnRemoved;
/// <summary>
/// Event for when the sync status changes
/// </summary>
event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// Start synchronization
/// </summary>
/// <returns></returns>
Task<CallResult> StartAsync(bool startWithSnapshot = true);
/// <summary>
/// Stop synchronization
/// </summary>
/// <returns></returns>
Task StopAsync();
/// <summary>
/// Get the data tracked
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
IEnumerable<SharedTrade> GetData(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
/// <summary>
/// Get statitistics on the trades
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
TradesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
}
}
@@ -0,0 +1,495 @@
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.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <inheritdoc />
public class TradeTracker : ITradeTracker
{
private readonly ITradeSocketClient _socketClient;
private readonly IRecentTradeRestClient? _recentRestClient;
private readonly ITradeHistoryRestClient? _historyRestClient;
private SyncStatus _status;
private long _snapshotId;
private bool _startWithSnapshot;
/// <summary>
/// The internal data structure
/// </summary>
protected readonly List<SharedTrade> _data = new List<SharedTrade>();
/// <summary>
/// The pre-snapshot queue buffering updates received before the snapshot is set and which will be applied after the snapshot was set
/// </summary>
protected readonly List<SharedTrade> _preSnapshotQueue = new List<SharedTrade>();
/// <summary>
/// The last time the window was applied
/// </summary>
protected DateTime _lastWindowApplied = DateTime.MinValue;
/// <summary>
/// Whether or not the data has changed since last window was applied
/// </summary>
protected bool _changed = false;
/// <summary>
/// Lock for accessing _data
/// </summary>
protected readonly object _lock = new object();
/// <summary>
/// Whether the snapshot has been set
/// </summary>
protected bool _snapshotSet;
/// <summary>
/// Logger
/// </summary>
protected readonly ILogger _logger;
/// <summary>
/// Update subscription
/// </summary>
protected UpdateSubscription? _updateSubscription;
/// <summary>
/// The timestamp of the first item
/// </summary>
protected DateTime? _firstTimestamp;
/// <inheritdoc />
public string Exchange { get; }
/// <inheritdoc />
public string SymbolName { get; }
/// <inheritdoc />
public SharedSymbol Symbol { get; }
/// <inheritdoc/>
public int? Limit { get; }
/// <inheritdoc/>
public TimeSpan? Period { get; }
/// <inheritdoc/>
public SyncStatus Status
{
get => _status;
set
{
if (value == _status)
return;
var old = _status;
_status = value;
_logger.TradeTrackerStatusChanged(SymbolName, old, value);
OnStatusChanged?.Invoke(old, _status);
}
}
/// <inheritdoc />
public int Count
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.Count;
}
}
}
/// <inheritdoc />
public DateTime? SyncedFrom
{
get
{
if (Period == null)
return _firstTimestamp;
var max = DateTime.UtcNow - Period.Value;
if (_firstTimestamp > max)
return _firstTimestamp;
return max;
}
}
/// <inheritdoc />
public SharedTrade? Last
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.LastOrDefault();
}
}
}
/// <inheritdoc />
public event Func<SharedTrade, Task>? OnAdded;
/// <inheritdoc />
public event Func<SharedTrade, Task>? OnRemoved;
/// <inheritdoc />
public event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// ctor
/// </summary>
public TradeTracker(
ILogger? logger,
IRecentTradeRestClient? recentRestClient,
ITradeHistoryRestClient? historyRestClient,
ITradeSocketClient socketClient,
SharedSymbol symbol,
int? limit = null,
TimeSpan? period = null)
{
_logger = logger ?? new NullLogger<TradeTracker>();
_recentRestClient = recentRestClient;
_historyRestClient = historyRestClient;
_socketClient = socketClient;
Exchange = socketClient.Exchange;
Symbol = symbol;
SymbolName = socketClient.FormatSymbol(symbol.BaseAsset, symbol.QuoteAsset, symbol.TradingMode, symbol.DeliverTime);
Limit = limit;
Period = period;
}
private TradesStats GetStats(IEnumerable<SharedTrade> trades)
{
if (!trades.Any())
return new TradesStats();
return new TradesStats
{
TradeCount = trades.Count(),
FirstTradeTime = trades.First().Timestamp,
LastTradeTime = trades.Last().Timestamp,
AveragePrice = Math.Round(trades.Select(d => d.Price).DefaultIfEmpty().Average(), 8),
VolumeWeightedAveragePrice = trades.Any() ? Math.Round(trades.Select(d => d.Price * d.Quantity).DefaultIfEmpty().Sum() / trades.Select(d => d.Quantity).DefaultIfEmpty().Sum(), 8) : null,
Volume = Math.Round(trades.Sum(d => d.Quantity), 8),
QuoteVolume = Math.Round(trades.Sum(d => d.Quantity * d.Price), 8),
BuySellRatio = Math.Round(trades.Where(x => x.Side == SharedOrderSide.Buy).Sum(x => x.Quantity) / trades.Sum(x => x.Quantity), 8)
};
}
/// <inheritdoc />
public TradesStats 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;
}
/// <inheritdoc />
public async Task<CallResult> 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.TradeTrackerStarting(SymbolName);
var subResult = await DoStartAsync().ConfigureAwait(false);
if (!subResult)
{
_logger.TradeTrackerStartFailed(SymbolName, subResult.Error!.ToString());
Status = SyncStatus.Disconnected;
return subResult;
}
_updateSubscription = subResult.Data;
_updateSubscription.ConnectionLost += HandleConnectionLost;
_updateSubscription.ConnectionClosed += HandleConnectionClosed;
_updateSubscription.ConnectionRestored += HandleConnectionRestored;
SetSyncStatus();
_logger.TradeTrackerStarted(SymbolName);
return new CallResult(null);
}
/// <inheritdoc />
public async Task StopAsync()
{
_logger.TradeTrackerStopping(SymbolName);
Status = SyncStatus.Disconnected;
await DoStopAsync().ConfigureAwait(false);
_data.Clear();
_preSnapshotQueue.Clear();
_logger.TradeTrackerStopped(SymbolName);
}
/// <summary>
/// The start procedure needed for trade syncing, generally subscribing to an update stream and requesting the snapshot
/// </summary>
/// <returns></returns>
protected virtual async Task<CallResult<UpdateSubscription>> DoStartAsync()
{
var subResult = await _socketClient.SubscribeToTradeUpdatesAsync(new SubscribeTradeRequest(Symbol),
update =>
{
AddData(update.Data);
}).ConfigureAwait(false);
if (!subResult)
{
Status = SyncStatus.Disconnected;
return subResult;
}
if (!_startWithSnapshot)
return subResult;
if (_historyRestClient != null)
{
var startTime = Period == null ? DateTime.UtcNow.AddMinutes(-5) : DateTime.UtcNow.Add(-Period.Value);
var request = new GetTradeHistoryRequest(Symbol, startTime, DateTime.UtcNow);
var data = new List<SharedTrade>();
await foreach(var result in ExchangeHelpers.ExecutePages(_historyRestClient.GetTradeHistoryAsync, request).ConfigureAwait(false))
{
if (!result)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(result.Error!);
}
if (Limit != null && data.Count > Limit)
break;
data.AddRange(result.Data);
}
SetInitialData(data);
}
else if (_recentRestClient != null)
{
int? limit = null;
if (Limit.HasValue)
limit = Math.Min(_recentRestClient.GetRecentTradesOptions.MaxLimit, Limit.Value);
var snapshot = await _recentRestClient.GetRecentTradesAsync(new GetRecentTradesRequest(Symbol, limit)).ConfigureAwait(false);
if (!snapshot)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(snapshot.Error!);
}
SetInitialData(snapshot.Data);
}
return subResult;
}
/// <summary>
/// The stop procedure needed, generally stopping the update stream
/// </summary>
/// <returns></returns>
protected virtual Task DoStopAsync() => _updateSubscription?.CloseAsync() ?? Task.CompletedTask;
/// <inheritdoc />
public IEnumerable<SharedTrade> GetData(DateTime? since = null, DateTime? until = null)
{
lock (_lock)
{
ApplyWindow(true);
IEnumerable<SharedTrade> result = _data;
if (since != null)
result = result.Where(d => d.Timestamp >= since);
if (until != null)
result = result.Where(d => d.Timestamp <= until);
return result.ToList();
}
}
/// <summary>
/// Set the initial trade data snapshot
/// </summary>
/// <param name="data"></param>
protected void SetInitialData(IEnumerable<SharedTrade> data)
{
lock (_lock)
{
_data.Clear();
IEnumerable<SharedTrade> items = data.OrderByDescending(d => d.Timestamp);
if (Limit != null)
items = items.Take(Limit.Value);
if (Period != null)
items = items.Where(e => e.Timestamp >= DateTime.UtcNow.Add(-Period.Value));
_snapshotId = data.Max(d => d.Timestamp.Ticks);
foreach (var item in items.OrderBy(d => d.Timestamp))
_data.Add(item);
_snapshotSet = true;
_changed = true;
_logger.TradeTrackerInitialDataSet(SymbolName, _data.Count, _snapshotId);
foreach (var item in _preSnapshotQueue)
{
if (_snapshotId >= item.Timestamp.Ticks)
{
_logger.TradeTrackerPreSnapshotSkip(SymbolName, item.Timestamp.Ticks);
continue;
}
_logger.TradeTrackerPreSnapshotApplied(SymbolName, item.Timestamp.Ticks);
_data.Add(item);
}
_firstTimestamp = _data.Min(v => v.Timestamp);
ApplyWindow(false);
}
}
/// <summary>
/// Add a trade
/// </summary>
/// <param name="item"></param>
protected void AddData(SharedTrade item) => AddData(new[] { item });
/// <summary>
/// Add a list of trades
/// </summary>
/// <param name="items"></param>
protected void AddData(IEnumerable<SharedTrade> items)
{
lock (_lock)
{
if ((_recentRestClient != null || _historyRestClient != null) && _startWithSnapshot && !_snapshotSet)
{
_preSnapshotQueue.AddRange(items);
return;
}
foreach (var item in items)
{
_logger.TradeTrackerTradeAdded(SymbolName, item.Timestamp.Ticks);
_data.Add(item);
OnAdded?.Invoke(item);
}
_firstTimestamp = _data.Min(x => x.Timestamp);
_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[0];
if (item.Timestamp >= compareDate)
break;
_data.Remove(item);
if (broadcastEvents)
OnRemoved?.Invoke(item);
}
}
if (Limit != null && _data.Count > Limit.Value)
{
var toRemove = _data.Count - Limit.Value;
for (var i = 0; i < toRemove; i++)
{
var item = _data[0];
_data.Remove(item);
if (broadcastEvents)
OnRemoved?.Invoke(item);
}
}
_lastWindowApplied = DateTime.UtcNow;
_changed = false;
if (Status == SyncStatus.PartiallySynced)
// Need to check if sync status should be changed even if there may not be any new data
SetSyncStatus();
}
private void HandleConnectionLost()
{
_logger.TradeTrackerConnectionLost(SymbolName);
if (Status != SyncStatus.Disconnected)
{
Status = SyncStatus.Syncing;
_snapshotSet = false;
_firstTimestamp = null;
_preSnapshotQueue.Clear();
}
}
private void HandleConnectionClosed()
{
_logger.TradeTrackerConnectionClosed(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;
}
_logger.TradeTrackerConnectionRestored(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;
}
}
}
@@ -0,0 +1,37 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// Trades statistics comparison
/// </summary>
public record TradesCompare
{
/// <summary>
/// Number of trades
/// </summary>
public CompareValue TradeCountDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// Average trade price
/// </summary>
public CompareValue? AveragePriceDif { get; set; }
/// <summary>
/// Volume weighted average trade price
/// </summary>
public CompareValue? VolumeWeightedAveragePriceDif { get; set; }
/// <summary>
/// Volume of the trades
/// </summary>
public CompareValue VolumeDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// Volume of the trades in quote asset
/// </summary>
public CompareValue QuoteVolumeDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// The volume weighted Buy/Sell ratio. A 0.7 ratio means 70% of the trade volume was a buy.
/// </summary>
public CompareValue? BuySellRatioDif { get; set; }
}
}
@@ -0,0 +1,65 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// Trades statistics
/// </summary>
public record TradesStats
{
/// <summary>
/// Number of trades
/// </summary>
public int TradeCount { get; set; }
/// <summary>
/// Timestamp of the last trade
/// </summary>
public DateTime? FirstTradeTime { get; set; }
/// <summary>
/// Timestamp of the first trade
/// </summary>
public DateTime? LastTradeTime { get; set; }
/// <summary>
/// Average trade price
/// </summary>
public decimal? AveragePrice { get; set; }
/// <summary>
/// Volume weighted average trade price
/// </summary>
public decimal? VolumeWeightedAveragePrice { get; set; }
/// <summary>
/// Volume of the trades
/// </summary>
public decimal Volume { get; set; }
/// <summary>
/// Volume of the trades in quote asset
/// </summary>
public decimal QuoteVolume { get; set; }
/// <summary>
/// The volume weighted Buy/Sell ratio. A 0.7 ratio means 70% of the trade volume was a buy.
/// </summary>
public decimal? BuySellRatio { get; set; }
/// <summary>
/// Whether the data is complete
/// </summary>
public bool Complete { get; set; }
/// <summary>
/// Compare 2 stat snapshots to eachother
/// </summary>
public TradesCompare CompareTo(TradesStats otherStats)
{
return new TradesCompare
{
TradeCountDif = new CompareValue(TradeCount, otherStats.TradeCount),
AveragePriceDif = new CompareValue(AveragePrice, otherStats.AveragePrice),
VolumeWeightedAveragePriceDif = new CompareValue(VolumeWeightedAveragePrice, otherStats.VolumeWeightedAveragePrice),
VolumeDif = new CompareValue(Volume, otherStats.Volume),
QuoteVolumeDif = new CompareValue(QuoteVolume, otherStats.QuoteVolume),
BuySellRatioDif = new CompareValue(BuySellRatio, otherStats.BuySellRatio),
};
}
}
}