1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-13 09:23:04 +00:00

Compare commits

...

14 Commits

17 changed files with 407 additions and 167 deletions
+3 -49
View File
@@ -304,55 +304,9 @@ namespace CryptoExchange.Net.Clients
return new CallResult<UpdateSubscription>(new ServerError(new ErrorInfo(ErrorType.WebsocketPaused, "Socket is paused"))); return new CallResult<UpdateSubscription>(new ServerError(new ErrorInfo(ErrorType.WebsocketPaused, "Socket is paused")));
} }
void HandleSubscriptionComplete(bool success, object? response) var subscribeResult = await socketConnection.TrySubscribeAsync(subscription, true, ct).ConfigureAwait(false);
{ if (!subscribeResult)
if (!success) return new CallResult<UpdateSubscription>(subscribeResult.Error!);
return;
subscription.HandleSubQueryResponse(socketConnection, response);
subscription.Status = SubscriptionStatus.Subscribed;
if (ct != default)
{
subscription.CancellationTokenRegistration = ct.Register(async () =>
{
_logger.CancellationTokenSetClosingSubscription(socketConnection.SocketId, subscription.Id);
await socketConnection.CloseAsync(subscription).ConfigureAwait(false);
}, false);
}
}
subscription.Status = SubscriptionStatus.Subscribing;
var subQuery = subscription.CreateSubscriptionQuery(socketConnection);
if (subQuery != null)
{
subQuery.OnComplete = () => HandleSubscriptionComplete(subQuery.Result?.Success ?? false, subQuery.Response);
// Send the request and wait for answer
var subResult = await socketConnection.SendAndWaitQueryAsync(subQuery, ct).ConfigureAwait(false);
if (!subResult)
{
var isTimeout = subResult.Error is CancellationRequestedError;
if (isTimeout && subscription.Status == SubscriptionStatus.Subscribed)
{
// No response received, but the subscription did receive updates. We'll assume success
}
else
{
_logger.FailedToSubscribe(socketConnection.SocketId, subResult.Error?.ToString());
// If this was a server process error we still might need to send an unsubscribe to prevent messages coming in later
subscription.Status = SubscriptionStatus.Pending;
await socketConnection.CloseAsync(subscription).ConfigureAwait(false);
return new CallResult<UpdateSubscription>(subResult.Error!);
}
}
if (!subQuery.ExpectsResponse)
HandleSubscriptionComplete(true, null);
}
else
{
HandleSubscriptionComplete(true, null);
}
_logger.SubscriptionCompletedSuccessfully(socketConnection.SocketId, subscription.Id); _logger.SubscriptionCompletedSuccessfully(socketConnection.SocketId, subscription.Id);
return new CallResult<UpdateSubscription>(new UpdateSubscription(socketConnection, subscription)); return new CallResult<UpdateSubscription>(new UpdateSubscription(socketConnection, subscription));
@@ -168,7 +168,7 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
if (!_unknownValuesWarned.Contains(stringValue)) if (!_unknownValuesWarned.Contains(stringValue))
{ {
_unknownValuesWarned.Add(stringValue!); _unknownValuesWarned.Add(stringValue!);
LibraryHelpers.StaticLogger?.LogWarning($"Cannot map enum value. EnumType: {enumType.FullName}, Value: {stringValue}, Known values: {string.Join(", ", _mappingToEnum!.Select(m => m.Value))}. If you think {stringValue} should added please open an issue on the Github repo"); LibraryHelpers.StaticLogger?.LogWarning($"Cannot map enum value. EnumType: {enumType.FullName}, Value: {stringValue}, Known values: [{string.Join(", ", _mappingToEnum!.Select(m => $"{m.StringValue}: {m.Value}"))}]. If you think {stringValue} should added please open an issue on the Github repo");
} }
} }
@@ -246,6 +246,12 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
{ {
// If no explicit mapping is found try to parse string // If no explicit mapping is found try to parse string
result = (T)Enum.Parse(objectType, value, true); result = (T)Enum.Parse(objectType, value, true);
if (!Enum.IsDefined(objectType, result))
{
result = default;
return false;
}
return true; return true;
} }
catch (Exception) catch (Exception)
+3 -3
View File
@@ -6,9 +6,9 @@
<PackageId>CryptoExchange.Net</PackageId> <PackageId>CryptoExchange.Net</PackageId>
<Authors>JKorf</Authors> <Authors>JKorf</Authors>
<Description>CryptoExchange.Net is a base library which is used to implement different cryptocurrency (exchange) API's. It provides a standardized way of implementing different API's, which results in a very similar experience for users of the API implementations.</Description> <Description>CryptoExchange.Net is a base library which is used to implement different cryptocurrency (exchange) API's. It provides a standardized way of implementing different API's, which results in a very similar experience for users of the API implementations.</Description>
<PackageVersion>10.5.0</PackageVersion> <PackageVersion>10.5.4</PackageVersion>
<AssemblyVersion>10.5.0</AssemblyVersion> <AssemblyVersion>10.5.4</AssemblyVersion>
<FileVersion>10.5.0</FileVersion> <FileVersion>10.5.4</FileVersion>
<PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance> <PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance>
<PackageTags>OKX;OKX.Net;Mexc;Mexc.Net;Kucoin;Kucoin.Net;Kraken;Kraken.Net;Huobi;Huobi.Net;CoinEx;CoinEx.Net;Bybit;Bybit.Net;Bitget;Bitget.Net;Bitfinex;Bitfinex.Net;Binance;Binance.Net;CryptoCurrency;CryptoCurrency Exchange;CryptoExchange.Net</PackageTags> <PackageTags>OKX;OKX.Net;Mexc;Mexc.Net;Kucoin;Kucoin.Net;Kraken;Kraken.Net;Huobi;Huobi.Net;CoinEx;CoinEx.Net;Bybit;Bybit.Net;Bitget;Bitget.Net;Bitfinex;Bitfinex.Net;Binance;Binance.Net;CryptoCurrency;CryptoCurrency Exchange;CryptoExchange.Net</PackageTags>
<RepositoryType>git</RepositoryType> <RepositoryType>git</RepositoryType>
@@ -1,5 +1,6 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Data.Common;
using System.Linq; using System.Linq;
namespace CryptoExchange.Net.SharedApis namespace CryptoExchange.Net.SharedApis
@@ -99,6 +100,9 @@ namespace CryptoExchange.Net.SharedApis
if (val == null) if (val == null)
return default; return default;
if (val.Value is T typeVal)
return typeVal;
try try
{ {
Type t = Nullable.GetUnderlyingType(typeof(T)) ?? typeof(T); Type t = Nullable.GetUnderlyingType(typeof(T)) ?? typeof(T);
@@ -673,13 +673,16 @@ namespace CryptoExchange.Net.Sockets.Default
if (!processed) if (!processed)
{ {
lock (_listenersLock) if (!ApiClient.HandleUnhandledMessage(this, typeIdentifier, data))
{ {
_logger.ReceivedMessageNotMatchedToAnyListener( lock (_listenersLock)
SocketId, {
typeIdentifier, _logger.ReceivedMessageNotMatchedToAnyListener(
topicFilter!, SocketId,
string.Join(",", _listeners.Select(x => string.Join(",", x.MessageRouter.Routes.Where(x => x.TypeIdentifier == typeIdentifier).Select(x => x.TopicFilter != null ? string.Join(",", x.TopicFilter) : "[null]"))))); typeIdentifier,
topicFilter!,
string.Join(",", _listeners.Select(x => string.Join(",", x.MessageRouter.Routes.Where(x => x.TypeIdentifier == typeIdentifier).Select(x => x.TopicFilter != null ? string.Join(",", x.TopicFilter) : "[null]")))));
}
} }
} }
} }
@@ -1082,40 +1085,8 @@ namespace CryptoExchange.Net.Sockets.Default
var taskList = new List<Task<CallResult>>(); var taskList = new List<Task<CallResult>>();
foreach (var subscription in subList) foreach (var subscription in subList)
{ {
subscription.ConnectionInvocations = 0; var subscribeTask = TrySubscribeAsync(subscription, false, default);
if (!subscription.Active) taskList.Add(subscribeTask);
// Can be closed during resubscribing
continue;
subscription.Status = SubscriptionStatus.Subscribing;
var result = await ApiClient.RevitalizeRequestAsync(subscription).ConfigureAwait(false);
if (!result)
{
_logger.FailedRequestRevitalization(SocketId, result.Error?.ToString());
subscription.Status = SubscriptionStatus.Pending;
return result;
}
var subQuery = subscription.CreateSubscriptionQuery(this);
if (subQuery == null)
{
subscription.Status = SubscriptionStatus.Subscribed;
continue;
}
subQuery.OnComplete = () =>
{
subscription.Status = subQuery.Result!.Success ? SubscriptionStatus.Subscribed : SubscriptionStatus.Pending;
subscription.HandleSubQueryResponse(this, subQuery.Response);
};
taskList.Add(SendAndWaitQueryAsync(subQuery));
if (!subQuery.ExpectsResponse)
{
// If there won't be an answer we can immediately set this
subscription.Status = SubscriptionStatus.Subscribed;
subscription.HandleSubQueryResponse(this, null);
}
} }
await Task.WhenAll(taskList).ConfigureAwait(false); await Task.WhenAll(taskList).ConfigureAwait(false);
@@ -1132,6 +1103,61 @@ namespace CryptoExchange.Net.Sockets.Default
return CallResult.SuccessResult; return CallResult.SuccessResult;
} }
protected internal async Task<CallResult> TrySubscribeAsync(Subscription subscription, bool newSubscription, CancellationToken subCancelToken)
{
subscription.ConnectionInvocations = 0;
if (!newSubscription)
{
if (!subscription.Active)
// Can be closed during resubscribing
return CallResult.SuccessResult;
var result = await ApiClient.RevitalizeRequestAsync(subscription).ConfigureAwait(false);
if (!result)
{
_logger.FailedRequestRevitalization(SocketId, result.Error?.ToString());
subscription.Status = SubscriptionStatus.Pending;
return result;
}
}
subscription.Status = SubscriptionStatus.Subscribing;
var subQuery = subscription.CreateSubscriptionQuery(this);
if (subQuery == null)
{
// No sub query, so successful
subscription.Status = SubscriptionStatus.Subscribed;
return CallResult.SuccessResult;
}
subQuery.OnComplete = () =>
{
subscription.Status = subQuery.Result!.Success ? SubscriptionStatus.Subscribed : SubscriptionStatus.Pending;
subscription.HandleSubQueryResponse(this, subQuery.Response);
if (newSubscription && subQuery.Result.Success && subCancelToken != default)
{
subscription.CancellationTokenRegistration = subCancelToken.Register(async () =>
{
_logger.CancellationTokenSetClosingSubscription(SocketId, subscription.Id);
await CloseAsync(subscription).ConfigureAwait(false);
}, false);
}
};
var subQueryResult = await SendAndWaitQueryAsync(subQuery).ConfigureAwait(false);
if (!subQueryResult)
{
_logger.FailedToSubscribe(SocketId, subQueryResult.Error?.ToString());
// If this was a server process error or timeout we still send an unsubscribe to prevent messages coming in later
if (newSubscription)
await CloseAsync(subscription).ConfigureAwait(false);
return new CallResult<UpdateSubscription>(subQueryResult.Error!);
}
return subQueryResult;
}
internal async Task UnsubscribeAsync(Subscription subscription) internal async Task UnsubscribeAsync(Subscription subscription)
{ {
var unsubscribeRequest = subscription.CreateUnsubscriptionQuery(this); var unsubscribeRequest = subscription.CreateUnsubscriptionQuery(this);
@@ -174,6 +174,14 @@ namespace CryptoExchange.Net.Sockets.Default
{ {
ConnectionInvocations++; ConnectionInvocations++;
TotalInvocations++; TotalInvocations++;
if (SubscriptionQuery != null && !SubscriptionQuery.Completed && SubscriptionQuery.TimeoutBehavior == TimeoutBehavior.Succeed)
{
// The subscription query is one where it is successful if there is no error returned
// Since we've received a data update for the subscription we can assume the subscribe query was successful
// Call timeout to complete
SubscriptionQuery.Timeout();
}
return route.Handle(connection, receiveTime, originalData, data); return route.Handle(connection, receiveTime, originalData, data);
} }
+3 -2
View File
@@ -127,8 +127,8 @@ namespace CryptoExchange.Net.Sockets
} }
else else
{ {
Completed = true;
Result = CallResult.SuccessResult; Result = CallResult.SuccessResult;
Completed = true;
_event.Set(); _event.Set();
} }
} }
@@ -216,12 +216,12 @@ namespace CryptoExchange.Net.Sockets
if (Completed) if (Completed)
return; return;
Completed = true;
if (TimeoutBehavior == TimeoutBehavior.Fail) if (TimeoutBehavior == TimeoutBehavior.Fail)
Result = new CallResult<THandlerResponse>(new TimeoutError()); Result = new CallResult<THandlerResponse>(new TimeoutError());
else else
Result = new CallResult<THandlerResponse>(default, null, default); Result = new CallResult<THandlerResponse>(default, null, default);
Completed = true;
_event.Set(); _event.Set();
OnComplete?.Invoke(); OnComplete?.Invoke();
} }
@@ -234,6 +234,7 @@ namespace CryptoExchange.Net.Sockets
Result = new CallResult<THandlerResponse>(error); Result = new CallResult<THandlerResponse>(error);
Completed = true; Completed = true;
_event.Set(); _event.Set();
OnComplete?.Invoke(); OnComplete?.Invoke();
} }
@@ -19,6 +19,7 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
private readonly IFuturesOrderSocketClient? _socketClient; private readonly IFuturesOrderSocketClient? _socketClient;
private readonly ExchangeParameters? _exchangeParameters; private readonly ExchangeParameters? _exchangeParameters;
private readonly bool _requiresSymbolParameterOpenOrders; private readonly bool _requiresSymbolParameterOpenOrders;
private readonly Dictionary<string, int> _openOrderNotReturnedTimes = new();
internal event Func<UpdateSource, SharedUserTrade[], Task>? OnTradeUpdate; internal event Func<UpdateSource, SharedUserTrade[], Task>? OnTradeUpdate;
@@ -251,11 +252,30 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
} }
} }
if (!_firstPollDone && anyError)
return anyError;
// Check all current open orders
// Keep track of the orders no longer returned in the open list
// Order should be set to canceled state when it's no longer returned in the open list
// but also is not returned in the closed list
foreach (var order in Values.Where(x => x.Status == SharedOrderStatus.Open))
{
if (openOrders.Any(x => x.OrderId == order.OrderId))
continue;
if (!_openOrderNotReturnedTimes.ContainsKey(order.OrderId))
_openOrderNotReturnedTimes[order.OrderId] = 0;
_openOrderNotReturnedTimes[order.OrderId] += 1;
}
var updatedPollTime = DateTime.UtcNow;
foreach (var symbol in _symbols.ToList()) foreach (var symbol in _symbols.ToList())
{ {
var fromTimeOrders = _lastDataTimeBeforeDisconnect ?? _lastPollTime ?? _startTime; DateTime? fromTimeOrders = GetClosedOrdersRequestStartTime(symbol);
var updatedPollTime = DateTime.UtcNow;
var closedOrdersResult = await _restClient.GetClosedFuturesOrdersAsync(new GetClosedOrdersRequest(symbol, startTime: fromTimeOrders, exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var closedOrdersResult = await _restClient.GetClosedFuturesOrdersAsync(new GetClosedOrdersRequest(symbol, startTime: fromTimeOrders, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!closedOrdersResult.Success) if (!closedOrdersResult.Success)
{ {
@@ -267,22 +287,26 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
else else
{ {
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
// Filter orders to only include where close time is after the start time // Filter orders to only include where close time is after the start time
var relevantOrders = closedOrdersResult.Data.Where(x => var relevantOrders = closedOrdersResult.Data.Where(x =>
(x.UpdateTime != null && x.UpdateTime >= _startTime) // Updated after the tracker start time (x.UpdateTime != null && x.UpdateTime >= _startTime) // Updated after the tracker start time
|| (x.CreateTime != null && x.CreateTime >= _startTime) // Created after the tracker start time || (x.CreateTime != null && x.CreateTime >= _startTime) // Created after the tracker start time
|| (x.CreateTime == null && x.UpdateTime == null) // Unknown time || (x.CreateTime == null && x.UpdateTime == null) // Unknown time
|| (Values.Any(e => e.OrderId == x.OrderId && x.Status == SharedOrderStatus.Open)) // Or we're currently tracking this open order
).ToArray(); ).ToArray();
// Check for orders which are no longer returned in either open/closed and assume they're canceled without fill // Check for orders which are no longer returned in either open/closed and assume they're canceled without fill
var openOrdersNotReturned = Values.Where(x => var openOrdersNotReturned = Values.Where(x =>
x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset // Orders for the same symbol // Orders for the same symbol
&& x.QuantityFilled?.IsZero == true // With no filled value x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset
&& !openOrders.Any(r => r.OrderId == x.OrderId) // Not returned in open orders // With no filled value
&& !relevantOrders.Any(r => r.OrderId == x.OrderId) // Not return in closed orders && x.QuantityFilled?.IsZero == true
// Not returned in open orders
&& !openOrders.Any(r => r.OrderId == x.OrderId)
// Not returned in closed orders
&& !relevantOrders.Any(r => r.OrderId == x.OrderId)
// Open order has not been returned in the open list at least 2 times
&& (_openOrderNotReturnedTimes.TryGetValue(x.OrderId, out var notReturnedTimes) ? notReturnedTimes >= 2 : false)
).ToList(); ).ToList();
var additionalUpdates = new List<SharedFuturesOrder>(); var additionalUpdates = new List<SharedFuturesOrder>();
@@ -300,7 +324,57 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
} }
if (!anyError)
{
_lastPollTime = updatedPollTime;
_lastDataTimeBeforeDisconnect = null;
}
return anyError; return anyError;
} }
private DateTime? GetClosedOrdersRequestStartTime(SharedSymbol symbol)
{
// Determine the timestamp from which we need to check order status
// Use the timestamp we last know the correct state of the data
DateTime? fromTime = null;
string? source = null;
// Use the last timestamp we we received data from the websocket as state should be correct at that time. 1 seconds buffer
if (_lastDataTimeBeforeDisconnect.HasValue && (fromTime == null || fromTime > _lastDataTimeBeforeDisconnect.Value))
{
fromTime = _lastDataTimeBeforeDisconnect.Value.AddSeconds(-1);
source = "LastDataTimeBeforeDisconnect";
}
// If we've previously polled use that timestamp to request data from
if (_lastPollTime.HasValue && (fromTime == null || _lastPollTime.Value > fromTime))
{
fromTime = _lastPollTime;
source = "LastPollTime";
}
// If we known open orders with a create time before this time we need to use that timestamp to make sure that order is included in the response
var trackedOrdersMinOpenTime = Values
.Where(x => x.Status == SharedOrderStatus.Open && x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset)
.OrderBy(x => x.CreateTime)
.FirstOrDefault()?.CreateTime;
if (trackedOrdersMinOpenTime.HasValue && (fromTime == null || trackedOrdersMinOpenTime.Value < fromTime))
{
// Could be improved by only requesting the specific open orders if there are only a few that would be better than trying to request a long
// history if the open order is far back
fromTime = trackedOrdersMinOpenTime.Value.AddMilliseconds(-1);
source = "OpenOrder";
}
if (fromTime == null)
{
fromTime = _startTime;
source = "StartTime";
}
_logger.LogTrace("{DataType} UserDataTracker poll startTime filter based on {Source}: {Time:yyyy-MM-dd HH:mm:ss.fff}", DataType, source, fromTime);
return fromTime!.Value;
}
} }
} }
@@ -55,10 +55,10 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
protected override async Task<bool> DoPollAsync() protected override async Task<bool> DoPollAsync()
{ {
var anyError = false; var anyError = false;
var fromTimeTrades = GetTradesRequestStartTime();
var updatedPollTime = DateTime.UtcNow;
foreach (var symbol in _symbols) foreach (var symbol in _symbols)
{ {
var fromTimeTrades = _lastDataTimeBeforeDisconnect ?? _lastPollTime ?? _startTime;
var updatedPollTime = DateTime.UtcNow;
var tradesResult = await _restClient.GetFuturesUserTradesAsync(new GetUserTradesRequest(symbol, startTime: fromTimeTrades, exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var tradesResult = await _restClient.GetFuturesUserTradesAsync(new GetUserTradesRequest(symbol, startTime: fromTimeTrades, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!tradesResult.Success) if (!tradesResult.Success)
{ {
@@ -80,9 +80,46 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
} }
if (!anyError)
{
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
}
return anyError; return anyError;
} }
private DateTime? GetTradesRequestStartTime()
{
// Determine the timestamp from which we need to check order status
// Use the timestamp we last know the correct state of the data
DateTime? fromTime = null;
string? source = null;
// Use the last timestamp we we received data from the websocket as state should be correct at that time. 1 seconds buffer
if (_lastDataTimeBeforeDisconnect.HasValue && (fromTime == null || fromTime > _lastDataTimeBeforeDisconnect.Value))
{
fromTime = _lastDataTimeBeforeDisconnect.Value.AddSeconds(-1);
source = "LastDataTimeBeforeDisconnect";
}
// If we've previously polled use that timestamp to request data from
if (_lastPollTime.HasValue && (fromTime == null || _lastPollTime.Value > fromTime))
{
fromTime = _lastPollTime;
source = "LastPollTime";
}
if (fromTime == null)
{
fromTime = _startTime;
source = "StartTime";
}
_logger.LogTrace("{DataType} UserDataTracker poll startTime filter based on {Source}: {Time:yyyy-MM-dd HH:mm:ss.fff}", DataType, source, fromTime);
return fromTime!.Value;
}
/// <inheritdoc /> /// <inheritdoc />
protected override Task<CallResult<UpdateSubscription?>> DoSubscribeAsync(string? listenKey) protected override Task<CallResult<UpdateSubscription?>> DoSubscribeAsync(string? listenKey)
{ {
@@ -19,6 +19,7 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
private readonly ISpotOrderSocketClient? _socketClient; private readonly ISpotOrderSocketClient? _socketClient;
private readonly ExchangeParameters? _exchangeParameters; private readonly ExchangeParameters? _exchangeParameters;
private readonly bool _requiresSymbolParameterOpenOrders; private readonly bool _requiresSymbolParameterOpenOrders;
private readonly Dictionary<string, int> _openOrderNotReturnedTimes = new();
internal event Func<UpdateSource, SharedUserTrade[], Task>? OnTradeUpdate; internal event Func<UpdateSource, SharedUserTrade[], Task>? OnTradeUpdate;
@@ -266,10 +267,27 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
if (!_firstPollDone && anyError) if (!_firstPollDone && anyError)
return anyError; return anyError;
// Check all current open orders
// Keep track of the orders no longer returned in the open list
// Order should be set to canceled state when it's no longer returned in the open list
// but also is not returned in the closed list
foreach (var order in Values.Where(x => x.Status == SharedOrderStatus.Open))
{
if (openOrders.Any(x => x.OrderId == order.OrderId))
continue;
if (!_openOrderNotReturnedTimes.ContainsKey(order.OrderId))
_openOrderNotReturnedTimes[order.OrderId] = 0;
_openOrderNotReturnedTimes[order.OrderId] += 1;
}
var updatedPollTime = DateTime.UtcNow;
foreach (var symbol in _symbols.ToList()) foreach (var symbol in _symbols.ToList())
{ {
var fromTimeOrders = _lastDataTimeBeforeDisconnect ?? _lastPollTime ?? _startTime; DateTime? fromTimeOrders = GetClosedOrdersRequestStartTime(symbol);
var updatedPollTime = DateTime.UtcNow;
var closedOrdersResult = await _restClient.GetClosedSpotOrdersAsync(new GetClosedOrdersRequest(symbol, startTime: fromTimeOrders, exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var closedOrdersResult = await _restClient.GetClosedSpotOrdersAsync(new GetClosedOrdersRequest(symbol, startTime: fromTimeOrders, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!closedOrdersResult.Success) if (!closedOrdersResult.Success)
{ {
@@ -281,22 +299,26 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
else else
{ {
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
// Filter orders to only include where close time is after the start time // Filter orders to only include where close time is after the start time
var relevantOrders = closedOrdersResult.Data.Where(x => var relevantOrders = closedOrdersResult.Data.Where(x =>
(x.UpdateTime != null && x.UpdateTime >= _startTime) // Updated after the tracker start time (x.UpdateTime != null && x.UpdateTime >= _startTime) // Updated after the tracker start time
|| (x.CreateTime != null && x.CreateTime >= _startTime) // Created after the tracker start time || (x.CreateTime != null && x.CreateTime >= _startTime) // Created after the tracker start time
|| (x.CreateTime == null && x.UpdateTime == null) // Unknown time || (x.CreateTime == null && x.UpdateTime == null) // Unknown time
|| (Values.Any(e => e.OrderId == x.OrderId && x.Status == SharedOrderStatus.Open)) // Or we're currently tracking this open order
).ToArray(); ).ToArray();
// Check for orders which are no longer returned in either open/closed and assume they're canceled without fill // Check for orders which are no longer returned in either open/closed and assume they're canceled without fill
var openOrdersNotReturned = Values.Where(x => var openOrdersNotReturned = Values.Where(x =>
x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset // Orders for the same symbol // Orders for the same symbol
&& x.QuantityFilled?.IsZero == true // With no filled value x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset
&& !openOrders.Any(r => r.OrderId == x.OrderId) // Not returned in open orders // With no filled value
&& !relevantOrders.Any(r => r.OrderId == x.OrderId) // Not return in closed orders && x.QuantityFilled?.IsZero == true
// Not returned in open orders
&& !openOrders.Any(r => r.OrderId == x.OrderId)
// Not returned in closed orders
&& !relevantOrders.Any(r => r.OrderId == x.OrderId)
// Open order has not been returned in the open list at least 2 times
&& (_openOrderNotReturnedTimes.TryGetValue(x.OrderId, out var notReturnedTimes) ? notReturnedTimes >= 2 : false)
).ToList(); ).ToList();
var additionalUpdates = new List<SharedSpotOrder>(); var additionalUpdates = new List<SharedSpotOrder>();
@@ -314,7 +336,57 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
} }
if (!anyError)
{
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
}
return anyError; return anyError;
} }
private DateTime? GetClosedOrdersRequestStartTime(SharedSymbol symbol)
{
// Determine the timestamp from which we need to check order status
// Use the timestamp we last know the correct state of the data
DateTime? fromTime = null;
string? source = null;
// Use the last timestamp we we received data from the websocket as state should be correct at that time. 1 seconds buffer
if (_lastDataTimeBeforeDisconnect.HasValue && (fromTime == null || fromTime > _lastDataTimeBeforeDisconnect.Value))
{
fromTime = _lastDataTimeBeforeDisconnect.Value.AddSeconds(-1);
source = "LastDataTimeBeforeDisconnect";
}
// If we've previously polled use that timestamp to request data from
if (_lastPollTime.HasValue && (fromTime == null || _lastPollTime.Value > fromTime))
{
fromTime = _lastPollTime;
source = "LastPollTime";
}
// If we known open orders with a create time before this time we need to use that timestamp to make sure that order is included in the response
var trackedOrdersMinOpenTime = Values
.Where(x => x.Status == SharedOrderStatus.Open && x.SharedSymbol!.BaseAsset == symbol.BaseAsset && x.SharedSymbol.QuoteAsset == symbol.QuoteAsset)
.OrderBy(x => x.CreateTime)
.FirstOrDefault()?.CreateTime;
if (trackedOrdersMinOpenTime.HasValue && (fromTime == null || trackedOrdersMinOpenTime.Value < fromTime))
{
// Could be improved by only requesting the specific open orders if there are only a few that would be better than trying to request a long
// history if the open order is far back
fromTime = trackedOrdersMinOpenTime.Value.AddMilliseconds(-1);
source = "OpenOrder";
}
if (fromTime == null)
{
fromTime = _startTime;
source = "StartTime";
}
_logger.LogTrace("{DataType} UserDataTracker poll startTime filter based on {Source}: {Time:yyyy-MM-dd HH:mm:ss.fff}", DataType, source, fromTime);
return fromTime!.Value;
}
} }
} }
@@ -55,10 +55,10 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
protected override async Task<bool> DoPollAsync() protected override async Task<bool> DoPollAsync()
{ {
var anyError = false; var anyError = false;
var fromTimeTrades = GetTradesRequestStartTime();
var updatedPollTime = DateTime.UtcNow;
foreach (var symbol in _symbols) foreach (var symbol in _symbols)
{ {
var fromTimeTrades = _lastDataTimeBeforeDisconnect ?? _lastPollTime ?? _startTime;
var updatedPollTime = DateTime.UtcNow;
var tradesResult = await _restClient.GetSpotUserTradesAsync(new GetUserTradesRequest(symbol, startTime: fromTimeTrades, exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var tradesResult = await _restClient.GetSpotUserTradesAsync(new GetUserTradesRequest(symbol, startTime: fromTimeTrades, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!tradesResult.Success) if (!tradesResult.Success)
{ {
@@ -70,8 +70,6 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
else else
{ {
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
// Filter trades to only include where timestamp is after the start time OR it's part of an order we're tracking // Filter trades to only include where timestamp is after the start time OR it's part of an order we're tracking
var relevantTrades = tradesResult.Data.Where(x => x.Timestamp >= _startTime || (GetTrackedOrderIds?.Invoke() ?? []).Any(o => o == x.OrderId)).ToArray(); var relevantTrades = tradesResult.Data.Where(x => x.Timestamp >= _startTime || (GetTrackedOrderIds?.Invoke() ?? []).Any(o => o == x.OrderId)).ToArray();
@@ -80,9 +78,46 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
} }
} }
if (!anyError)
{
_lastDataTimeBeforeDisconnect = null;
_lastPollTime = updatedPollTime;
}
return anyError; return anyError;
} }
private DateTime? GetTradesRequestStartTime()
{
// Determine the timestamp from which we need to check order status
// Use the timestamp we last know the correct state of the data
DateTime? fromTime = null;
string? source = null;
// Use the last timestamp we we received data from the websocket as state should be correct at that time. 1 seconds buffer
if (_lastDataTimeBeforeDisconnect.HasValue && (fromTime == null || fromTime > _lastDataTimeBeforeDisconnect.Value))
{
fromTime = _lastDataTimeBeforeDisconnect.Value.AddSeconds(-1);
source = "LastDataTimeBeforeDisconnect";
}
// If we've previously polled use that timestamp to request data from
if (_lastPollTime.HasValue && (fromTime == null || _lastPollTime.Value > fromTime))
{
fromTime = _lastPollTime;
source = "LastPollTime";
}
if (fromTime == null)
{
fromTime = _startTime;
source = "StartTime";
}
_logger.LogTrace("{DataType} UserDataTracker poll startTime filter based on {Source}: {Time:yyyy-MM-dd HH:mm:ss.fff}", DataType, source, fromTime);
return fromTime!.Value;
}
/// <inheritdoc /> /// <inheritdoc />
protected override Task<CallResult<UpdateSubscription?>> DoSubscribeAsync(string? listenKey) protected override Task<CallResult<UpdateSubscription?>> DoSubscribeAsync(string? listenKey)
{ {
@@ -372,6 +372,7 @@ namespace CryptoExchange.Net.Trackers.UserData.ItemTrackers
{ {
toRemove ??= new List<T>(); toRemove ??= new List<T>();
toRemove.Add(item); toRemove.Add(item);
_logger.LogWarning("Ignoring {DataType} update for {Key}, no SharedSymbol set", DataType, GetKey(item));
} }
else if (_onlyTrackProvidedSymbols else if (_onlyTrackProvidedSymbols
&& !_symbols.Any(y => y.TradingMode == symbolModel.SharedSymbol!.TradingMode && y.BaseAsset == symbolModel.SharedSymbol.BaseAsset && y.QuoteAsset == symbolModel.SharedSymbol.QuoteAsset)) && !_symbols.Any(y => y.TradingMode == symbolModel.SharedSymbol!.TradingMode && y.BaseAsset == symbolModel.SharedSymbol.BaseAsset && y.QuoteAsset == symbolModel.SharedSymbol.QuoteAsset))
@@ -20,6 +20,7 @@ namespace CryptoExchange.Net.Trackers.UserData
private readonly IFuturesSymbolRestClient _symbolClient; private readonly IFuturesSymbolRestClient _symbolClient;
private readonly IListenKeyRestClient? _listenKeyClient; private readonly IListenKeyRestClient? _listenKeyClient;
private readonly ExchangeParameters? _exchangeParameters; private readonly ExchangeParameters? _exchangeParameters;
private readonly TradingMode _tradingMode;
private Task? _lkKeepAliveTask; private Task? _lkKeepAliveTask;
/// <inheritdoc /> /// <inheritdoc />
@@ -69,10 +70,14 @@ namespace CryptoExchange.Net.Trackers.UserData
_listenKeyClient = listenKeyRestClient; _listenKeyClient = listenKeyRestClient;
_exchangeParameters = exchangeParameters; _exchangeParameters = exchangeParameters;
_tradingMode = accountType == SharedAccountType.PerpetualInverseFutures ? TradingMode.PerpetualInverse :
accountType == SharedAccountType.DeliveryLinearFutures ? TradingMode.DeliveryLinear :
accountType == SharedAccountType.DeliveryInverseFutures ? TradingMode.DeliveryInverse :
TradingMode.PerpetualLinear;
var trackers = new List<UserDataItemTracker>(); var trackers = new List<UserDataItemTracker>();
var balanceAccountType = accountType ?? SharedAccountType.PerpetualLinearFutures; var balanceTracker = new BalanceTracker(logger, balanceRestClient, balanceSocketClient, accountType ?? SharedAccountType.PerpetualLinearFutures, config.BalancesConfig, exchangeParameters);
var balanceTracker = new BalanceTracker(logger, balanceRestClient, balanceSocketClient, balanceAccountType, config.BalancesConfig, exchangeParameters);
Balances = balanceTracker; Balances = balanceTracker;
trackers.Add(balanceTracker); trackers.Add(balanceTracker);
@@ -100,7 +105,7 @@ namespace CryptoExchange.Net.Trackers.UserData
/// <inheritdoc /> /// <inheritdoc />
protected override async Task<CallResult> DoStartAsync() protected override async Task<CallResult> DoStartAsync()
{ {
var symbolResult = await _symbolClient.GetFuturesSymbolsAsync(new GetSymbolsRequest(exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var symbolResult = await _symbolClient.GetFuturesSymbolsAsync(new GetSymbolsRequest(_tradingMode, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!symbolResult) if (!symbolResult)
{ {
_logger.LogWarning("Failed to start UserFuturesDataTracker; symbols request failed: {Error}", symbolResult.Error); _logger.LogWarning("Failed to start UserFuturesDataTracker; symbols request failed: {Error}", symbolResult.Error);
@@ -109,7 +114,7 @@ namespace CryptoExchange.Net.Trackers.UserData
if (_listenKeyClient != null) if (_listenKeyClient != null)
{ {
var lkResult = await _listenKeyClient.StartListenKeyAsync(new StartListenKeyRequest(exchangeParameters: _exchangeParameters)).ConfigureAwait(false); var lkResult = await _listenKeyClient.StartListenKeyAsync(new StartListenKeyRequest(_tradingMode, exchangeParameters: _exchangeParameters)).ConfigureAwait(false);
if (!lkResult) if (!lkResult)
{ {
_logger.LogWarning("Failed to start UserFuturesDataTracker; listen key request failed: {Error}", lkResult.Error); _logger.LogWarning("Failed to start UserFuturesDataTracker; listen key request failed: {Error}", lkResult.Error);
@@ -141,7 +146,7 @@ namespace CryptoExchange.Net.Trackers.UserData
break; break;
} }
var result = await _listenKeyClient!.KeepAliveListenKeyAsync(new KeepAliveListenKeyRequest(_listenKey!, TradingMode.Spot)).ConfigureAwait(false); var result = await _listenKeyClient!.KeepAliveListenKeyAsync(new KeepAliveListenKeyRequest(_listenKey!, _tradingMode)).ConfigureAwait(false);
if (!result) if (!result)
_logger.LogWarning("Listen key keep alive failed: " + result.Error); _logger.LogWarning("Listen key keep alive failed: " + result.Error);
+25 -25
View File
@@ -5,32 +5,32 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="12.1.0" /> <PackageReference Include="Binance.Net" Version="12.5.0" />
<PackageReference Include="Bitfinex.Net" Version="10.2.0" /> <PackageReference Include="Bitfinex.Net" Version="10.6.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" /> <PackageReference Include="BitMart.Net" Version="3.5.0" />
<PackageReference Include="BloFin.Net" Version="2.1.1" /> <PackageReference Include="BloFin.Net" Version="2.5.0" />
<PackageReference Include="Bybit.Net" Version="6.1.0" /> <PackageReference Include="Bybit.Net" Version="6.5.0" />
<PackageReference Include="CoinEx.Net" Version="10.1.0" /> <PackageReference Include="CoinEx.Net" Version="10.5.0" />
<PackageReference Include="CoinW.Net" Version="2.1.1" /> <PackageReference Include="CoinW.Net" Version="2.5.0" />
<PackageReference Include="CryptoCom.Net" Version="3.1.0" /> <PackageReference Include="CryptoCom.Net" Version="3.5.0" />
<PackageReference Include="DeepCoin.Net" Version="3.1.0" /> <PackageReference Include="DeepCoin.Net" Version="3.5.0" />
<PackageReference Include="GateIo.Net" Version="3.1.0" /> <PackageReference Include="GateIo.Net" Version="3.5.0" />
<PackageReference Include="HyperLiquid.Net" Version="3.2.0" /> <PackageReference Include="HyperLiquid.Net" Version="3.7.0" />
<PackageReference Include="JK.BingX.Net" Version="3.1.0" /> <PackageReference Include="JK.BingX.Net" Version="3.5.0" />
<PackageReference Include="JK.Bitget.Net" Version="3.1.0" /> <PackageReference Include="JK.Bitget.Net" Version="3.5.0" />
<PackageReference Include="JK.Mexc.Net" Version="4.1.0" /> <PackageReference Include="JK.Mexc.Net" Version="4.5.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" /> <PackageReference Include="JK.OKX.Net" Version="4.5.0" />
<PackageReference Include="Jkorf.Aster.Net" Version="2.1.0" /> <PackageReference Include="Jkorf.Aster.Net" Version="2.5.0" />
<PackageReference Include="JKorf.BitMEX.Net" Version="3.1.0" /> <PackageReference Include="JKorf.BitMEX.Net" Version="3.5.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="3.1.0" /> <PackageReference Include="JKorf.Coinbase.Net" Version="3.5.1" />
<PackageReference Include="JKorf.HTX.Net" Version="8.1.0" /> <PackageReference Include="JKorf.HTX.Net" Version="8.5.0" />
<PackageReference Include="JKorf.Upbit.Net" Version="2.1.0" /> <PackageReference Include="JKorf.Upbit.Net" Version="2.5.0" />
<PackageReference Include="KrakenExchange.Net" Version="7.1.0" /> <PackageReference Include="KrakenExchange.Net" Version="7.5.0" />
<PackageReference Include="Kucoin.Net" Version="8.1.0" /> <PackageReference Include="Kucoin.Net" Version="8.5.0" />
<PackageReference Include="Serilog.AspNetCore" Version="10.0.0" /> <PackageReference Include="Serilog.AspNetCore" Version="10.0.0" />
<PackageReference Include="Toobit.Net" Version="2.1.0" /> <PackageReference Include="Toobit.Net" Version="3.5.0" />
<PackageReference Include="WhiteBit.Net" Version="3.1.0" /> <PackageReference Include="WhiteBit.Net" Version="3.5.0" />
<PackageReference Include="XT.Net" Version="3.1.0" /> <PackageReference Include="XT.Net" Version="3.5.0" />
</ItemGroup> </ItemGroup>
</Project> </Project>
+14 -14
View File
@@ -6,20 +6,20 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="12.1.0" /> <PackageReference Include="Binance.Net" Version="12.5.0" />
<PackageReference Include="Bitfinex.Net" Version="10.2.0" /> <PackageReference Include="Bitfinex.Net" Version="10.6.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" /> <PackageReference Include="BitMart.Net" Version="3.5.0" />
<PackageReference Include="Bybit.Net" Version="6.1.0" /> <PackageReference Include="Bybit.Net" Version="6.5.0" />
<PackageReference Include="CoinEx.Net" Version="10.1.0" /> <PackageReference Include="CoinEx.Net" Version="10.5.0" />
<PackageReference Include="CryptoCom.Net" Version="3.1.0" /> <PackageReference Include="CryptoCom.Net" Version="3.5.0" />
<PackageReference Include="GateIo.Net" Version="3.1.0" /> <PackageReference Include="GateIo.Net" Version="3.5.0" />
<PackageReference Include="JK.Bitget.Net" Version="3.1.0" /> <PackageReference Include="JK.Bitget.Net" Version="3.5.0" />
<PackageReference Include="JK.Mexc.Net" Version="4.1.0" /> <PackageReference Include="JK.Mexc.Net" Version="4.5.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" /> <PackageReference Include="JK.OKX.Net" Version="4.5.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="3.1.0" /> <PackageReference Include="JKorf.Coinbase.Net" Version="3.5.1" />
<PackageReference Include="JKorf.HTX.Net" Version="8.1.0" /> <PackageReference Include="JKorf.HTX.Net" Version="8.5.0" />
<PackageReference Include="KrakenExchange.Net" Version="7.1.0" /> <PackageReference Include="KrakenExchange.Net" Version="7.5.0" />
<PackageReference Include="Kucoin.Net" Version="8.1.0" /> <PackageReference Include="Kucoin.Net" Version="8.5.0" />
</ItemGroup> </ItemGroup>
</Project> </Project>
+3 -3
View File
@@ -8,9 +8,9 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="12.1.0" /> <PackageReference Include="Binance.Net" Version="12.5.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" /> <PackageReference Include="BitMart.Net" Version="3.5.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" /> <PackageReference Include="JK.OKX.Net" Version="4.5.0" />
</ItemGroup> </ItemGroup>
</Project> </Project>
+17
View File
@@ -67,6 +67,23 @@ Make a one time donation in a crypto currency of your choice. If you prefer to d
Alternatively, sponsor me on Github using [Github Sponsors](https://github.com/sponsors/JKorf). Alternatively, sponsor me on Github using [Github Sponsors](https://github.com/sponsors/JKorf).
## Release notes ## Release notes
* Version 10.5.4 - 12 Feb 2026
* Fixed type check ExchangeParameters GetValue
* Fixed bug in polling time filter for UserDataTracker item
* Version 10.5.3 - 11 Feb 2026
* Fixed orders getting incorrectly set to canceled state for UserDataTracker spot and futures orders
* Added check EnumConverter to detect undefined int value parsing
* Version 10.5.2 - 10 Feb 2026
* Added check for subscribe queries with TimeoutBehavior.Success to complete when subscription has received update
* Added call to ApiClient.HandleUnhandledMessage when no websocket message processor is found based on topic to allow additional processing
* Combined websocket connection subscribe and re-subscribe logic
* Set websocket query completed after setting Result
* Version 10.5.1 - 10 Feb 2026
* Fixed trading mode selection for futures listen key methods in FuturesUserDataTracker
* Version 10.5.0 - 10 Feb 2026 * Version 10.5.0 - 10 Feb 2026
* Added keep alive for listenkeys to UserDataTracker * Added keep alive for listenkeys to UserDataTracker
* Updated logging unmatched websocket message * Updated logging unmatched websocket message