1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-12 17:03:10 +00:00

Compare commits

..

28 Commits

Author SHA1 Message Date
JKorf 5942423bfb Updated to version 10.2.4 2026-01-17 16:29:56 +01:00
JKorf dc4abc42a7 Added WaitUntilFirstUpdateBufferedAsync method on SymbolOrderBook, fixed sequencen validation bug SymbolOrderBook 2026-01-17 16:26:25 +01:00
Jkorf c71a81e686 Added some util methods, Added CommaSplitStringConverter 2026-01-16 16:39:21 +01:00
Jkorf 550c0eabf1 Updated to version 10.2.3 2026-01-14 08:46:07 +01:00
JKorf 28a2a0c7fd Fixed semaphore exception when creating a new REST client while time sync is in progress on another client 2026-01-13 22:11:18 +01:00
JKorf 7dd1cd5bbd Added HandleUnhandledMessage virtual method to SocketApiClient to allow some processing for messages which couldn't be mapped via the normal way 2026-01-13 21:21:30 +01:00
Jkorf a7ff4416bd Updated to version 10.2.2 2026-01-13 11:18:16 +01:00
Jkorf 669d1f7c9e Allow the same websocket connection sequence number to be recorded multiple times 2026-01-13 11:16:37 +01:00
Jkorf 34ee2d3690 Updated to version 10.2.1 2026-01-13 09:29:04 +01:00
Jkorf 005fb7875d Fixed parameter URL creation for array values with ArrayParametersSerialization.MultipleValues 2026-01-13 09:25:35 +01:00
Jkorf fa9300ce97 Removed duplicate logging for rest responses in Trace verbosity 2026-01-13 09:25:07 +01:00
Jkorf fc2d3fc2d2 Updated to version 10.2.0 2026-01-12 14:30:38 +01:00
Jkorf 187ca6a4ef Fixed warning 2026-01-12 14:30:09 +01:00
Jan Korf 3b2a85d210 Feature/websocket sequencing (#267)
Added EnforceSequenceNumbers property on SocketApiClient to configure whether websocket message contain sequence numbers and if these should be checked to be sequential
Added fallback to existing websocket connection if no dedicated request connection was found
Added IntBoolConverter base class for arbitrary int value to bool mapping
Added SequenceNumber property to DataEvent object
Added _skipSequenceCheckFirstUpdateAfterSnapshotSet property for SymbolOrderBook implementations
Updated SymbolOrderBook sequenceNumber validation
Updated SymbolOrderBook log verbosities
Renamed SetInitialOrderBook to SetSnapshot in SymbolOrderBook
Renamed updateId references to sequenceNumber in SymbolOrderBook
2026-01-12 14:26:50 +01:00
Jkorf c512bee825 Updated examples 2026-01-07 14:56:18 +01:00
Jkorf 0943b052b9 Updated to version 10.1.0 2026-01-07 10:03:45 +01:00
Jan Korf a896fffdb3 Time offset management (#266)
Updated time sync / time offset management for REST API's
Added time offset tracking for WebSocket API's
Added GetAuthenticationQuery virtual method on AuthenticationProvider
Updated AuthenticationProvider GetTimestamp methods to include a one second offset by default
Added AuthenticationProvider GetTimestamp methods for SocketApiClient instances
Added ClientName property on BaseApiClient, resolving to the type name
Added ObjectOrArrayConverter JsonConverterFactory implementation for resolving json data which might be returned as object or array
Added UpdateServerTime, UpdateLocalTime and DataAge properties to (I)SymbolOrderBook
Added OutputToConsoleAsync method to (I)SymbolOrderBook
Updated SymbolOrderBook string representation
Added DataTimeLocal and DataAge properties to DataEvent object
Added SocketConnection parameter to subscription HandleSubQueryResponse and HandleUnsubQueryResponse methods
2026-01-07 10:00:14 +01:00
JKorf 177daf903b Added some utils methods 2025-12-30 09:53:30 +01:00
Jkorf aa1ebdc4ed Updated CryptoExchange.Net to version 10.0.2 2025-12-19 11:34:26 +01:00
Jkorf 38058c4a70 Updated to version 10.0.2 2025-12-19 10:14:59 +01:00
Jkorf a7eb483479 Added exception handlers for REST response processing 2025-12-19 10:09:14 +01:00
Jkorf c76931a3b4 Fixed duplicate subscription check with updated deserialization 2025-12-19 09:58:22 +01:00
Jkorf b90b7e9e0c Updated CryptoExchange.Net to 10.0.1 2025-12-18 11:13:56 +01:00
Jkorf beda53d36d Updated to version 10.0.1 2025-12-18 10:56:10 +01:00
Jkorf 0668f669c1 Fixed query parameter array serialization 2025-12-18 10:36:22 +01:00
Jkorf 64250e13db Updated examples 2025-12-17 10:55:48 +01:00
Jkorf 451d38d5e7 Updated to version 10.0.1 2025-12-16 11:54:37 +01:00
Jkorf e11e437bbb Fixed CryptoExchange.Net reference 2025-12-16 11:54:08 +01:00
32 changed files with 1095 additions and 354 deletions
@@ -6,9 +6,9 @@
<PackageId>CryptoExchange.Net.Protobuf</PackageId>
<Authors>JKorf</Authors>
<Description>Protobuf support for CryptoExchange.Net</Description>
<PackageVersion>10.0.0</PackageVersion>
<AssemblyVersion>10.0.0</AssemblyVersion>
<FileVersion>10.0.0</FileVersion>
<PackageVersion>10.0.1</PackageVersion>
<AssemblyVersion>10.0.1</AssemblyVersion>
<FileVersion>10.0.1</FileVersion>
<PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance>
<PackageTags>CryptoExchange;CryptoExchange.Net</PackageTags>
<RepositoryType>git</RepositoryType>
@@ -41,7 +41,7 @@
<DocumentationFile>CryptoExchange.Net.Protobuf.xml</DocumentationFile>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="CryptoExchange.Net" Version="9.13.0" />
<PackageReference Include="CryptoExchange.Net" Version="10.0.2" />
<PackageReference Include="protobuf-net" Version="3.2.56" />
</ItemGroup>
</Project>
+3
View File
@@ -5,6 +5,9 @@
Protobuf support for CryptoExchange.Net.
## Release notes
* Version 10.0.1 - 16 Dec 2025
* Updated CryptoExchange.Net version to 10.0.0, see https://github.com/JKorf/CryptoExchange.Net/releases/
* Version 10.0.0 - 16 Dec 2025
* Updated CryptoExchange.Net version to 10.0.0, see https://github.com/JKorf/CryptoExchange.Net/releases/
@@ -9,7 +9,7 @@
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="18.0.1"></PackageReference>
<PackageReference Include="Moq" Version="4.20.72" />
<PackageReference Include="NUnit" Version="4.4.0"></PackageReference>
<PackageReference Include="NUnit3TestAdapter" Version="6.0.0"></PackageReference>
<PackageReference Include="NUnit3TestAdapter" Version="6.0.1"></PackageReference>
</ItemGroup>
<ItemGroup>
@@ -63,8 +63,6 @@ namespace CryptoExchange.Net.UnitTests
/// <inheritdoc />
public override string FormatSymbol(string baseAsset, string quoteAsset, TradingMode futuresType, DateTime? deliverDate = null) => $"{baseAsset.ToUpperInvariant()}{quoteAsset.ToUpperInvariant()}";
public override TimeSpan? GetTimeOffset() => null;
public override TimeSyncInfo GetTimeSyncInfo() => null;
protected override IStreamMessageAccessor CreateAccessor() => new SystemTextJsonStreamMessageAccessor(new System.Text.Json.JsonSerializerOptions());
protected override IMessageSerializer CreateSerializer() => new SystemTextJsonMessageSerializer(new System.Text.Json.JsonSerializerOptions());
protected override AuthenticationProvider CreateAuthenticationProvider(ApiCredentials credentials) => throw new NotImplementedException();
@@ -160,11 +160,6 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
ParameterPositions[method] = position;
}
public override TimeSpan? GetTimeOffset()
{
throw new NotImplementedException();
}
protected override AuthenticationProvider CreateAuthenticationProvider(ApiCredentials credentials)
=> new TestAuthProvider(credentials);
@@ -172,11 +167,6 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
{
throw new NotImplementedException();
}
public override TimeSyncInfo GetTimeSyncInfo()
{
throw new NotImplementedException();
}
}
public class TestRestApi2Client : RestApiClient
@@ -198,12 +188,7 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
{
return await SendAsync<T>("http://www.test.com", new RequestDefinition("/", HttpMethod.Get) { Weight = 0 }, null, ct);
}
public override TimeSpan? GetTimeOffset()
{
throw new NotImplementedException();
}
protected override AuthenticationProvider CreateAuthenticationProvider(ApiCredentials credentials)
=> new TestAuthProvider(credentials);
@@ -212,10 +197,6 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
throw new NotImplementedException();
}
public override TimeSyncInfo GetTimeSyncInfo()
{
throw new NotImplementedException();
}
}
public class TestError
@@ -11,6 +11,8 @@ using System.Linq;
using System.Globalization;
using System.Security.Cryptography;
using System.Text;
using CryptoExchange.Net.Sockets;
using CryptoExchange.Net.Sockets.Default;
namespace CryptoExchange.Net.Authentication
{
@@ -76,12 +78,20 @@ namespace CryptoExchange.Net.Authentication
}
/// <summary>
/// Authenticate a request
/// Authenticate a REST request
/// </summary>
/// <param name="apiClient">The Api client sending the request</param>
/// <param name="apiClient">The API client sending the request</param>
/// <param name="requestConfig">The request configuration</param>
public abstract void ProcessRequest(RestApiClient apiClient, RestRequestConfiguration requestConfig);
/// <summary>
/// Get an authentication query for a websocket
/// </summary>
/// <param name="apiClient">The API client sending the request</param>
/// <param name="connection">The connection to authenticate</param>
/// <param name="context">Optional context required for creating the authentication query</param>
public virtual Query? GetAuthenticationQuery(SocketApiClient apiClient, SocketConnection connection, Dictionary<string, object?>? context = null) => null;
/// <summary>
/// SHA256 sign the data and return the bytes
/// </summary>
@@ -442,6 +452,14 @@ namespace CryptoExchange.Net.Authentication
/// <param name="buff"></param>
/// <returns></returns>
protected static string BytesToHexString(byte[] buff)
=> BytesToHexString(new ArraySegment<byte>(buff));
/// <summary>
/// Convert byte array to hex string
/// </summary>
/// <param name="buff"></param>
/// <returns></returns>
protected static string BytesToHexString(ArraySegment<byte> buff)
{
#if NET9_0_OR_GREATER
return Convert.ToHexString(buff);
@@ -453,6 +471,26 @@ namespace CryptoExchange.Net.Authentication
#endif
}
/// <summary>
/// Convert a hex encoded string to byte array
/// </summary>
/// <param name="hexString"></param>
/// <returns></returns>
protected static byte[] HexToBytesString(string hexString)
{
if (hexString.StartsWith("0x"))
hexString = hexString.Substring(2);
byte[] bytes = new byte[hexString.Length / 2];
for (int i = 0; i < hexString.Length; i += 2)
{
string hexSubstring = hexString.Substring(i, 2);
bytes[i / 2] = Convert.ToByte(hexSubstring, 16);
}
return bytes;
}
/// <summary>
/// Convert byte array to base64 string
/// </summary>
@@ -466,32 +504,50 @@ namespace CryptoExchange.Net.Authentication
/// <summary>
/// Get current timestamp including the time sync offset from the api client
/// </summary>
/// <param name="apiClient"></param>
/// <returns></returns>
protected DateTime GetTimestamp(RestApiClient apiClient)
protected DateTime GetTimestamp(RestApiClient apiClient, bool includeOneSecondOffset = true)
{
return TimeProvider.GetTime().Add(apiClient.GetTimeOffset() ?? TimeSpan.Zero)!;
var result = TimeProvider.GetTime().Add(TimeOffsetManager.GetRestOffset(apiClient.ClientName) ?? TimeSpan.Zero)!;
if (includeOneSecondOffset)
result = result.AddSeconds(-1);
return result;
}
/// <summary>
/// Get current timestamp including the time sync offset from the api client
/// </summary>
protected DateTime GetTimestamp(SocketApiClient apiClient, bool includeOneSecondOffset = true)
{
var result = TimeProvider.GetTime().Add(TimeOffsetManager.GetSocketOffset(apiClient.ClientName) ?? TimeSpan.Zero)!;
if (includeOneSecondOffset)
result = result.AddSeconds(-1);
return result;
}
/// <summary>
/// Get millisecond timestamp as a string including the time sync offset from the api client
/// </summary>
/// <param name="apiClient"></param>
/// <returns></returns>
protected string GetMillisecondTimestamp(RestApiClient apiClient)
{
return DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient)).Value.ToString(CultureInfo.InvariantCulture);
}
protected string GetMillisecondTimestamp(RestApiClient apiClient, bool includeOneSecondOffset = true)
=> DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient, includeOneSecondOffset)).Value.ToString(CultureInfo.InvariantCulture);
/// <summary>
/// Get millisecond timestamp as a string including the time sync offset from the api client
/// </summary>
protected string GetMillisecondTimestamp(SocketApiClient apiClient, bool includeOneSecondOffset = true)
=> DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient, includeOneSecondOffset)).Value.ToString(CultureInfo.InvariantCulture);
/// <summary>
/// Get millisecond timestamp as a long including the time sync offset from the api client
/// </summary>
/// <param name="apiClient"></param>
/// <returns></returns>
protected long GetMillisecondTimestampLong(RestApiClient apiClient)
{
return DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient)).Value;
}
protected long GetMillisecondTimestampLong(RestApiClient apiClient, bool includeOneSecondOffset = true)
=> DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient, includeOneSecondOffset)).Value;
/// <summary>
/// Get millisecond timestamp as a long including the time sync offset from the api client
/// </summary>
protected long GetMillisecondTimestampLong(SocketApiClient apiClient, bool includeOneSecondOffset = true)
=> DateTimeConverter.ConvertToMilliseconds(GetTimestamp(apiClient, includeOneSecondOffset)).Value;
/// <summary>
/// Return the serialized request body
@@ -13,6 +13,8 @@ namespace CryptoExchange.Net.Clients
/// </summary>
public abstract class BaseApiClient : IDisposable, IBaseApiClient
{
private string? _clientName;
/// <summary>
/// Logger
/// </summary>
@@ -23,6 +25,21 @@ namespace CryptoExchange.Net.Clients
/// </summary>
protected bool _disposing;
/// <summary>
/// Name of the client
/// </summary>
protected internal string ClientName
{
get
{
if (_clientName != null)
return _clientName;
_clientName = GetType().Name;
return _clientName;
}
}
/// <summary>
/// The authentication provider for this API client. (null if no credentials are set)
/// </summary>
+76 -59
View File
@@ -32,12 +32,6 @@ namespace CryptoExchange.Net.Clients
/// <inheritdoc />
public IRequestFactory RequestFactory { get; set; } = new RequestFactory();
/// <inheritdoc />
public abstract TimeSyncInfo? GetTimeSyncInfo();
/// <inheritdoc />
public abstract TimeSpan? GetTimeOffset();
/// <inheritdoc />
public int TotalRequestsMade { get; set; }
@@ -115,6 +109,8 @@ namespace CryptoExchange.Net.Clients
options,
apiOptions)
{
TimeOffsetManager.RegisterRestApi(ClientName);
RequestFactory.Configure(options, httpClient);
}
@@ -241,11 +237,9 @@ namespace CryptoExchange.Net.Clients
{
currentTry++;
var error = await CheckTimeSync(requestId, definition).ConfigureAwait(false);
if (error != null)
return new WebCallResult<T>(error);
await CheckTimeSync(requestId, definition).ConfigureAwait(false);
error = await RateLimitAsync(
var error = await RateLimitAsync(
baseAddress,
requestId,
definition,
@@ -300,28 +294,6 @@ namespace CryptoExchange.Net.Clients
}
}
private async ValueTask<Error?> CheckTimeSync(int requestId, RequestDefinition definition)
{
if (!definition.Authenticated)
return null;
var syncTask = SyncTimeAsync();
var timeSyncInfo = GetTimeSyncInfo();
if (timeSyncInfo != null && timeSyncInfo.TimeSyncState.LastSyncTime == default)
{
// Initially with first request we'll need to wait for the time syncing, if it's not the first request we can just continue
var syncTimeError = await syncTask.ConfigureAwait(false);
if (syncTimeError != null)
{
_logger.RestApiFailedToSyncTime(requestId, syncTimeError!.ToString());
return syncTimeError;
}
}
return null;
}
/// <summary>
/// Check rate limits for the request
/// </summary>
@@ -482,9 +454,6 @@ namespace CryptoExchange.Net.Clients
{
memoryStream.Position = 0;
originalData = await reader.ReadToEndAsync().ConfigureAwait(false);
if (_logger.IsEnabled(LogLevel.Trace))
_logger.RestApiReceivedResponse(request.RequestId, originalData);
}
// Continue processing from the memory stream since the response stream is already read and we can't seek it
@@ -516,10 +485,20 @@ namespace CryptoExchange.Net.Clients
else
{
// Handle a 'normal' error response. Can still be either a json error message or some random HTML or other string
error = await MessageHandler.ParseErrorResponse(
try
{
error = await MessageHandler.ParseErrorResponse(
(int)response.StatusCode,
response.ResponseHeaders,
responseStream).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "Unhandled exception when parsing error response: {Message}", ex.Message);
var errorResult = new ServerError(ErrorInfo.Unknown with { Message = ex.Message });
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, default, errorResult);
}
}
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, default, error);
@@ -558,10 +537,19 @@ namespace CryptoExchange.Net.Clients
if (deserializeError != null)
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, deserializeResult, deserializeError); ;
// Check the deserialized response to see if it's an error or not
var responseError = MessageHandler.CheckDeserializedResponse(response.ResponseHeaders, deserializeResult);
if (responseError != null)
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, deserializeResult, responseError);
try
{
// Check the deserialized response to see if it's an error or not
var responseError = MessageHandler.CheckDeserializedResponse(response.ResponseHeaders, deserializeResult);
if (responseError != null)
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, deserializeResult, responseError);
}
catch (Exception ex)
{
_logger.LogError(ex, "Unhandled exception when checking deserialized response: {Message}", ex.Message);
var error = new ServerError(ErrorInfo.Unknown with { Message = ex.Message });
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, deserializeResult, error);
}
return new WebCallResult<T>(response.StatusCode, response.HttpVersion, response.ResponseHeaders, sw.Elapsed, response.ContentLength, originalData, request.RequestId, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), ResultDataSource.Server, deserializeResult, null);
}
@@ -706,26 +694,44 @@ namespace CryptoExchange.Net.Clients
RequestFactory.UpdateSettings(options.Proxy, options.RequestTimeout ?? ClientOptions.RequestTimeout, ClientOptions.HttpKeepAliveInterval);
}
internal async ValueTask<Error?> SyncTimeAsync()
private async ValueTask CheckTimeSync(int requestId, RequestDefinition definition)
{
var timeSyncParams = GetTimeSyncInfo();
if (timeSyncParams == null)
return null;
if (!definition.Authenticated)
return;
if (await timeSyncParams.TimeSyncState.Semaphore.WaitAsync(0).ConfigureAwait(false))
var lastUpdateTime = TimeOffsetManager.GetRestLastUpdateTime(ClientName);
var syncTask = CheckTimeOffsetAsync();
if (lastUpdateTime == null)
{
if (!timeSyncParams.SyncTime || DateTime.UtcNow - timeSyncParams.TimeSyncState.LastSyncTime < timeSyncParams.RecalculationInterval)
{
timeSyncParams.TimeSyncState.Semaphore.Release();
return null;
}
// Initially with first request we'll need to wait for the time syncing before making the actual request.
// If it's not the first request we can just continue and let it complete in the background
await syncTask.ConfigureAwait(false);
}
return;
}
internal async ValueTask CheckTimeOffsetAsync()
{
if (!(ApiOptions.AutoTimestamp ?? ClientOptions.AutoTimestamp))
// Time syncing not enabled
return;
await TimeOffsetManager.EnterAsync(ClientName).ConfigureAwait(false);
try
{
var lastUpdateTime = TimeOffsetManager.GetRestLastUpdateTime(ClientName);
if (DateTime.UtcNow - lastUpdateTime < (ApiOptions.TimestampRecalculationInterval ?? ClientOptions.TimestampRecalculationInterval))
// Time syncing was recently done
return;
var localTime = DateTime.UtcNow;
var result = await GetServerTimestampAsync().ConfigureAwait(false);
if (!result)
{
timeSyncParams.TimeSyncState.Semaphore.Release();
return result.Error;
_logger.LogWarning("Failed to determine time offset between client and server, timestamping might fail");
return;
}
if (TotalRequestsMade == 1)
@@ -735,18 +741,29 @@ namespace CryptoExchange.Net.Clients
result = await GetServerTimestampAsync().ConfigureAwait(false);
if (!result)
{
timeSyncParams.TimeSyncState.Semaphore.Release();
return result.Error;
_logger.LogWarning("Failed to determine time offset between client and server, timestamping might fail");
return;
}
}
// Calculate time offset between local and server
// Estimate the offset as the round trip time / 2
var offset = result.Data - localTime.AddMilliseconds(result.ResponseTime!.Value.TotalMilliseconds / 2);
timeSyncParams.UpdateTimeOffset(offset);
timeSyncParams.TimeSyncState.Semaphore.Release();
}
if (offset.TotalMilliseconds > 0 && offset.TotalMilliseconds < 500)
{
_logger.LogInformation("{ClientName} Time offset within limits ({Offset}ms), set offset to 0ms", ClientName, Math.Round(offset.TotalMilliseconds));
offset = TimeSpan.Zero;
}
else
{
_logger.LogInformation("{ClientName} Time offset set to {Offset}ms", ClientName, Math.Round(offset.TotalMilliseconds));
}
return null;
TimeOffsetManager.UpdateRestOffset(ClientName, offset.TotalMilliseconds);
}
finally
{
TimeOffsetManager.Release(ClientName);
}
}
private bool ShouldCache(RequestDefinition definition)
+39 -3
View File
@@ -32,8 +32,10 @@ namespace CryptoExchange.Net.Clients
public abstract class SocketApiClient : BaseApiClient, ISocketApiClient
{
#region Fields
/// <inheritdoc/>
public IWebsocketFactory SocketFactory { get; set; } = new WebsocketFactory();
/// <inheritdoc/>
public IHighPerfConnectionFactory? HighPerfConnectionFactory { get; set; }
@@ -140,6 +142,10 @@ namespace CryptoExchange.Net.Clients
/// </summary>
public int? MaxIndividualSubscriptionsPerConnection { get; set; }
/// <summary>
/// Whether or not to enforce that sequence number updates are always (lastSequenceNumber + 1)
/// </summary>
public bool EnforceSequenceNumbers { get; set; }
#endregion
/// <summary>
@@ -181,6 +187,24 @@ namespace CryptoExchange.Net.Clients
DedicatedConnectionConfigs.Add(new DedicatedConnectionConfig() { SocketAddress = url, Authenticated = auth });
}
/// <summary>
/// Update the timestamp offset between client and server based on the timestamp
/// </summary>
/// <param name="timestamp">Timestamp received from the server</param>
public virtual void UpdateTimeOffset(DateTime timestamp)
{
if (timestamp == default)
return;
TimeOffsetManager.UpdateSocketOffset(ClientName, (DateTime.UtcNow - timestamp).TotalMilliseconds);
}
/// <summary>
/// Get the time offset between client and server
/// </summary>
/// <returns></returns>
public virtual TimeSpan? GetTimeOffset() => TimeOffsetManager.GetSocketOffset(ClientName);
/// <summary>
/// Add a query to periodically send on each connection
/// </summary>
@@ -296,7 +320,7 @@ namespace CryptoExchange.Net.Clients
if (!success)
return;
subscription.HandleSubQueryResponse(response);
subscription.HandleSubQueryResponse(socketConnection, response);
subscription.Status = SubscriptionStatus.Subscribed;
if (ct != default)
{
@@ -575,7 +599,8 @@ namespace CryptoExchange.Net.Clients
/// Should return the request which can be used to authenticate a socket connection
/// </summary>
/// <returns></returns>
protected internal virtual Task<Query?> GetAuthenticationRequestAsync(SocketConnection connection) => throw new NotImplementedException();
protected internal virtual Task<Query?> GetAuthenticationRequestAsync(SocketConnection connection) =>
Task.FromResult(AuthenticationProvider!.GetAuthenticationQuery(this, connection));
/// <summary>
/// Adds a system subscription. Used for example to reply to ping requests
@@ -685,6 +710,10 @@ namespace CryptoExchange.Net.Clients
if (connection != null && !connection.DedicatedRequestConnection.Authenticated)
// Mark dedicated request connection as authenticated if the request is authenticated
connection.DedicatedRequestConnection.Authenticated = authenticated;
if (connection == null)
// Fall back to an existing connection if there is no dedicated request connection available
connection = socketQuery.OrderBy(s => s.UserSubscriptionCount).FirstOrDefault();
}
bool maxConnectionsReached = _socketConnections.Count >= (ApiOptions.MaxSocketConnections ?? ClientOptions.MaxSocketConnections);
@@ -773,7 +802,6 @@ namespace CryptoExchange.Net.Clients
return new CallResult<HighPerfSocketConnection<TUpdateType>>(socketConnection);
}
/// <summary>
/// Process an unhandled message
/// </summary>
@@ -782,6 +810,14 @@ namespace CryptoExchange.Net.Clients
{
}
/// <summary>
/// Process an unhandled message
/// </summary>
/// <param name="connection">The socket connection</param>
/// <param name="typeIdentifier">The type as identified</param>
/// <param name="data">The data</param>
protected internal virtual bool HandleUnhandledMessage(SocketConnection connection, string typeIdentifier, ReadOnlySpan<byte> data) => false;
/// <summary>
/// Process connect rate limited
/// </summary>
@@ -0,0 +1,30 @@
using System;
using System.Diagnostics.CodeAnalysis;
using System.Linq;
using System.Text.Json;
using System.Text.Json.Serialization;
namespace CryptoExchange.Net.Converters.SystemTextJson
{
/// <summary>
/// Converter for comma separated string values
/// </summary>
public class CommaSplitStringConverter : JsonConverter<string[]>
{
/// <inheritdoc />
public override string[]? Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
{
var str = reader.GetString();
if (string.IsNullOrEmpty(str))
return [];
return str!.Split(',').ToArray() ?? [];
}
/// <inheritdoc />
public override void Write(Utf8JsonWriter writer, string[] value, JsonSerializerOptions options)
{
writer.WriteStringValue(string.Join(",", value));
}
}
}
@@ -0,0 +1,38 @@
using System;
using System.Text.Json;
using System.Text.Json.Serialization;
namespace CryptoExchange.Net.Converters.SystemTextJson
{
/// <summary>
/// Bool converter
/// </summary>
public class IntBoolConverter : JsonConverter<bool>
{
private readonly int _trueValue;
/// <summary>
/// ctor
/// </summary>
/// <param name="trueValue">The int value representing the true value</param>
public IntBoolConverter(int trueValue)
{
_trueValue = trueValue;
}
/// <inheritdoc />
public override bool Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
{
if (reader.TokenType != JsonTokenType.Number)
return false;
return reader.GetDecimal() == _trueValue;
}
/// <inheritdoc />
public override void Write(Utf8JsonWriter writer, bool value, JsonSerializerOptions options)
{
writer.WriteNumberValue(_trueValue);
}
}
}
@@ -0,0 +1,58 @@
using System;
using System.Text.Json;
using System.Text.Json.Serialization;
#pragma warning disable IL2026 // Members annotated with 'RequiresUnreferencedCodeAttribute' require dynamic access otherwise can break functionality when trimming application code
#pragma warning disable IL3050 // Calling members annotated with 'RequiresDynamicCodeAttribute' may break functionality when AOT compiling.
namespace CryptoExchange.Net.Converters.SystemTextJson
{
/// <summary>
/// Converter for parsing object or array responses
/// </summary>
public class ObjectOrArrayConverter : JsonConverterFactory
{
/// <inheritdoc />
public override bool CanConvert(Type typeToConvert) => true;
/// <inheritdoc />
public override JsonConverter? CreateConverter(Type typeToConvert, JsonSerializerOptions options)
{
var type = typeof(InternalObjectOrArrayConverter<>).MakeGenericType(typeToConvert);
return (JsonConverter)Activator.CreateInstance(type)!;
}
private class InternalObjectOrArrayConverter<T> : JsonConverter<T>
{
public override T? Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
{
if (reader.TokenType == JsonTokenType.StartObject && !typeToConvert.IsArray)
{
// Object to object
return JsonDocument.ParseValue(ref reader).Deserialize<T>(options);
}
else if (reader.TokenType == JsonTokenType.StartArray && typeToConvert.IsArray)
{
// Array to array
return JsonDocument.ParseValue(ref reader).Deserialize<T>(options);
}
else if (reader.TokenType == JsonTokenType.StartArray)
{
// Array to object
JsonDocument.ParseValue(ref reader).Deserialize<T[]>(options);
return default;
}
else
{
// Object to array
JsonDocument.ParseValue(ref reader);
return default;
}
}
public override void Write(Utf8JsonWriter writer, T value, JsonSerializerOptions options)
{
JsonSerializer.Serialize(writer, value, options);
}
}
}
}
+3 -3
View File
@@ -6,9 +6,9 @@
<PackageId>CryptoExchange.Net</PackageId>
<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>
<PackageVersion>10.0.0</PackageVersion>
<AssemblyVersion>10.0.0</AssemblyVersion>
<FileVersion>10.0.0</FileVersion>
<PackageVersion>10.2.4</PackageVersion>
<AssemblyVersion>10.2.4</AssemblyVersion>
<FileVersion>10.2.4</FileVersion>
<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>
<RepositoryType>git</RepositoryType>
+20 -2
View File
@@ -242,8 +242,7 @@ namespace CryptoExchange.Net
/// <summary>
/// Generate a long value
/// </summary>
/// <param name="maxLength">Max character length</param>
/// <returns></returns>
/// <param name="maxLength">Max number of digits</param>
public static long RandomLong(int maxLength)
{
#if NETSTANDARD2_1_OR_GREATER || NET9_0_OR_GREATER
@@ -259,6 +258,25 @@ namespace CryptoExchange.Net
return value;
}
/// <summary>
/// Generate a long value between two values
/// </summary>
/// <param name="minValue">Min value</param>
/// <param name="maxValue">Max value</param>
/// <returns></returns>
public static long RandomLong(long minValue, long maxValue)
{
#if NET8_0_OR_GREATER
var buf = RandomNumberGenerator.GetBytes(8);
#else
byte[] buf = new byte[8];
var random = new Random();
random.NextBytes(buf);
#endif
long longRand = BitConverter.ToInt64(buf, 0);
return (Math.Abs(longRand % (maxValue - minValue)) + minValue);
}
/// <summary>
/// Generate a random string of specified length
/// </summary>
+31 -2
View File
@@ -70,12 +70,17 @@ namespace CryptoExchange.Net
first = false;
if (parameter.GetType().IsArray)
if (parameter.Value.GetType().IsArray)
{
if (serializationType == ArrayParametersSerialization.Array)
{
foreach(var entry in (object[])parameter.Value)
bool firstArrayValue = true;
foreach (var entry in (object[])parameter.Value)
{
if (!firstArrayValue)
uriString.Append('&');
firstArrayValue = false;
uriString.Append(parameter.Key);
uriString.Append("[]=");
if (urlEncodeValues)
@@ -86,8 +91,12 @@ namespace CryptoExchange.Net
}
else if (serializationType == ArrayParametersSerialization.MultipleValues)
{
bool firstArrayValue = true;
foreach (var entry in (object[])parameter.Value)
{
if (!firstArrayValue)
uriString.Append('&');
firstArrayValue = false;
uriString.Append(parameter.Key);
uriString.Append("=");
if (urlEncodeValues)
@@ -607,6 +616,26 @@ namespace CryptoExchange.Net
return services;
}
/// <summary>
/// Convert a hex encoded string to byte array
/// </summary>
/// <param name="hexString"></param>
/// <returns></returns>
public static byte[] HexStringToBytes(this string hexString)
{
if (hexString.StartsWith("0x"))
hexString = hexString.Substring(2);
byte[] bytes = new byte[hexString.Length / 2];
for (int i = 0; i < hexString.Length; i += 2)
{
string hexSubstring = hexString.Substring(i, 2);
bytes[i / 2] = Convert.ToByte(hexSubstring, 16);
}
return bytes;
}
}
}
@@ -47,9 +47,21 @@ namespace CryptoExchange.Net.Interfaces
/// </summary>
event Action<(ISymbolOrderBookEntry BestBid, ISymbolOrderBookEntry BestAsk)> OnBestOffersChanged;
/// <summary>
/// Timestamp of the last update
/// Timestamp of when the last update was applied to the book, local time
/// </summary>
DateTime UpdateTime { get; }
/// <summary>
/// Timestamp of the last event that was applied, server time
/// </summary>
DateTime? UpdateServerTime { get; }
/// <summary>
/// Timestamp of the last event that was applied, in local time, estimated based on timestamp difference between client and server
/// </summary>
DateTime? UpdateLocalTime { get; }
/// <summary>
/// Age of the data, in local time, estimated based on timestamp difference between client and server + the period since last update
/// </summary>
TimeSpan? DataAge { get; }
/// <summary>
/// The number of asks in the book
@@ -126,5 +138,13 @@ namespace CryptoExchange.Net.Interfaces
/// </summary>
/// <returns></returns>
string ToString(int rows);
/// <summary>
/// Output the orderbook to the console
/// </summary>
/// <param name="numberOfEntries">Number of rows to display</param>
/// <param name="refreshInterval">Refresh interval</param>
/// <param name="ct">Cancellation token</param>
Task OutputToConsoleAsync(int numberOfEntries, TimeSpan refreshInterval, CancellationToken ct = default);
}
}
@@ -22,7 +22,6 @@ namespace CryptoExchange.Net.Logging.Extensions
private static readonly Action<ILogger, string, Exception?> _restApiCacheHit;
private static readonly Action<ILogger, string, Exception?> _restApiCacheNotHit;
private static readonly Action<ILogger, int?, Exception?> _restApiCancellationRequested;
private static readonly Action<ILogger, int?, string?, Exception?> _restApiReceivedResponse;
static RestApiClientLoggingExtensions()
{
@@ -90,11 +89,6 @@ namespace CryptoExchange.Net.Logging.Extensions
LogLevel.Debug,
new EventId(4012, "RestApiCancellationRequested"),
"[Req {RequestId}] Request cancelled by user");
_restApiReceivedResponse = LoggerMessage.Define<int?, string?>(
LogLevel.Trace,
new EventId(4013, "RestApiReceivedResponse"),
"[Req {RequestId}] Received response: {Data}");
}
@@ -161,10 +155,5 @@ namespace CryptoExchange.Net.Logging.Extensions
{
_restApiCancellationRequested(logger, requestId, null);
}
public static void RestApiReceivedResponse(this ILogger logger, int requestId, string? originalData)
{
_restApiReceivedResponse(logger, requestId, originalData, null);
}
}
}
@@ -22,13 +22,15 @@ namespace CryptoExchange.Net.Logging.Extensions
private static readonly Action<ILogger, string, string, Exception?> _orderBookResyncing;
private static readonly Action<ILogger, string, string, Exception?> _orderBookResynced;
private static readonly Action<ILogger, string, string, Exception?> _orderBookMessageSkippedBecauseOfResubscribing;
private static readonly Action<ILogger, string, string, long, long, long, Exception?> _orderBookDataSet;
private static readonly Action<ILogger, string, string, long, long, long?, Exception?> _orderBookDataSet;
private static readonly Action<ILogger, string, string, long, long, long, long, Exception?> _orderBookUpdateBuffered;
private static readonly Action<ILogger, string, string, decimal, decimal, Exception?> _orderBookOutOfSyncDetected;
private static readonly Action<ILogger, string, string, Exception?> _orderBookReconnectingSocket;
private static readonly Action<ILogger, string, string, long, long, Exception?> _orderBookSkippedMessage;
private static readonly Action<ILogger, string, string, long, long, Exception?> _orderBookProcessedMessage;
private static readonly Action<ILogger, string, string, long, Exception?> _orderBookProcessedMessageSingle;
private static readonly Action<ILogger, string, string, long, long, Exception?> _orderBookOutOfSync;
private static readonly Action<ILogger, string, string, long, long, long, Exception?> _orderBookUpdateSkippedStartEnd;
static SymbolOrderBookLoggingExtensions()
{
@@ -73,7 +75,7 @@ namespace CryptoExchange.Net.Logging.Extensions
"{Api} order book {Symbol} Processing {NumberBufferedUpdated} buffered updates");
_orderBookUpdateSkipped = LoggerMessage.Define<string, string, long, long>(
LogLevel.Debug,
LogLevel.Trace,
new EventId(5008, "OrderBookUpdateSkipped"),
"{Api} order book {Symbol} update skipped #{SequenceNumber}, currently at #{LastSequenceNumber}");
@@ -92,10 +94,10 @@ namespace CryptoExchange.Net.Logging.Extensions
new EventId(5011, "OrderBookMessageSkippedResubscribing"),
"{Api} order book {Symbol} Skipping message because of resubscribing");
_orderBookDataSet = LoggerMessage.Define<string, string, long, long, long>(
LogLevel.Debug,
_orderBookDataSet = LoggerMessage.Define<string, string, long, long, long?>(
LogLevel.Trace,
new EventId(5012, "OrderBookDataSet"),
"{Api} order book {Symbol} data set: {BidCount} bids, {AskCount} asks. #{EndUpdateId}");
"{Api} order book {Symbol} snapshot set: {BidCount} bids, {AskCount} asks. #{EndUpdateId}");
_orderBookUpdateBuffered = LoggerMessage.Define<string, string, long, long, long, long>(
LogLevel.Trace,
@@ -136,6 +138,17 @@ namespace CryptoExchange.Net.Logging.Extensions
LogLevel.Warning,
new EventId(5020, "OrderBookOutOfSyncChecksum"),
"{Api} order book {Symbol} out of sync. Checksum mismatch, resyncing");
_orderBookProcessedMessageSingle = LoggerMessage.Define<string, string, long>(
LogLevel.Trace,
new EventId(5021, "OrderBookProcessedMessage"),
"{Api} order book {Symbol} update processed #{UpdateId}");
_orderBookUpdateSkippedStartEnd = LoggerMessage.Define<string, string, long, long, long>(
LogLevel.Trace,
new EventId(5022, "OrderBookUpdateSkippedStartEnd"),
"{Api} order book {Symbol} update skipped #{SequenceStart}-#{SequenceEnd}, currently at #{LastSequenceNumber}");
}
public static void OrderBookStatusChanged(this ILogger logger, string api, string symbol, OrderBookStatus previousStatus, OrderBookStatus newStatus)
@@ -194,7 +207,7 @@ namespace CryptoExchange.Net.Logging.Extensions
{
_orderBookMessageSkippedBecauseOfResubscribing(logger, api, symbol, null);
}
public static void OrderBookDataSet(this ILogger logger, string api, string symbol, long bidCount, long askCount, long endUpdateId)
public static void OrderBookDataSet(this ILogger logger, string api, string symbol, long bidCount, long askCount, long? endUpdateId)
{
_orderBookDataSet(logger, api, symbol, bidCount, askCount, endUpdateId, null);
}
@@ -229,9 +242,18 @@ namespace CryptoExchange.Net.Logging.Extensions
_orderBookProcessedMessage(logger, api, symbol, firstUpdateId, lastUpdateId, null);
}
public static void OrderBookProcessedMessage(this ILogger logger, string api, string symbol, long updateId)
{
_orderBookProcessedMessageSingle(logger, api, symbol, updateId, null);
}
public static void OrderBookOutOfSyncChecksum(this ILogger logger, string api, string symbol)
{
_orderBookOutOfSyncChecksum(logger, api, symbol, null);
}
public static void OrderBookUpdateSkipped(this ILogger logger, string api, string symbol, long sequenceStart, long sequenceEnd, long lastSequenceNumber)
{
_orderBookUpdateSkippedStartEnd(logger, api, symbol, sequenceStart, sequenceEnd, lastSequenceNumber, null);
}
}
}
@@ -96,6 +96,23 @@ namespace CryptoExchange.Net.Objects
base.Add(key, value.Value.ToString(CultureInfo.InvariantCulture));
}
/// <summary>
/// Add a DateTime value as string
/// </summary>
public void AddString(string key, DateTime value)
{
base.Add(key, value.ToString("yyyy-MM-ddTHH:mm:ssZ"));
}
/// <summary>
/// Add a DateTime value as string. Not added if value is null
/// </summary>
public void AddOptionalString(string key, DateTime? value)
{
if (value != null)
base.Add(key, value.Value.ToString("yyyy-MM-ddTHH:mm:ssZ"));
}
/// <summary>
/// Add a datetime value as milliseconds timestamp
/// </summary>
@@ -241,6 +258,45 @@ namespace CryptoExchange.Net.Objects
base.Add(key, int.Parse(stringVal));
}
}
/// <summary>
/// Add key as comma separated values
/// </summary>
public void AddCommaSeparated(string key, IEnumerable<string> values)
{
base.Add(key, string.Join(",", values));
}
/// <summary>
/// Add key as comma separated values if there are values provided
/// </summary>
public void AddOptionalCommaSeparated(string key, IEnumerable<string>? values)
{
if (values == null || !values.Any())
return;
base.Add(key, string.Join(",", values));
}
/// <summary>
/// Add key as boolean lower case value
/// </summary>
public void AddBoolString(string key, bool value)
{
base.Add(key, value.ToString().ToLower());
}
/// <summary>
/// Add key as boolean lower case value if it's not null
/// </summary>
public void AddOptionalBoolString(string key, bool? value)
{
if (value == null)
return;
base.Add(key, value.ToString()!.ToLower());
}
/// <summary>
/// Set the request body. Can be used to specify a simple value or array as the body instead of an object
@@ -18,6 +18,16 @@ namespace CryptoExchange.Net.Objects.Sockets
/// </summary>
public DateTime? DataTime { get; set; }
/// <summary>
/// The timestamp of the data in local time. Note that this is an estimation based on average delay from the server.
/// </summary>
public DateTime? DataTimeLocal { get; set; }
/// <summary>
/// The age of the data. Note that this is an estimation based on average delay from the server.
/// </summary>
public TimeSpan? DataAge => DateTime.UtcNow - DataTimeLocal;
/// <summary>
/// The stream producing the update
/// </summary>
@@ -43,6 +53,11 @@ namespace CryptoExchange.Net.Objects.Sockets
/// </summary>
public SocketUpdateType? UpdateType { get; set; }
/// <summary>
/// Sequence number of the update
/// </summary>
public long? SequenceNumber { get; set; }
/// <summary>
/// ctor
/// </summary>
@@ -116,12 +131,28 @@ namespace CryptoExchange.Net.Objects.Sockets
return this;
}
/// <summary>
/// Specify the sequence number of the update
/// </summary>
public DataEvent<T> WithSequenceNumber(long? sequenceNumber)
{
SequenceNumber = sequenceNumber;
return this;
}
/// <summary>
/// Specify the data timestamp
/// </summary>
public DataEvent<T> WithDataTimestamp(DateTime? timestamp)
public DataEvent<T> WithDataTimestamp(DateTime? timestamp, TimeSpan? offset)
{
if (timestamp == null || timestamp == default(DateTime))
return this;
DataTime = timestamp;
if (offset == null)
return this;
DataTimeLocal = DataTime + offset;
return this;
}
@@ -1,95 +0,0 @@
using System;
using System.Threading;
using Microsoft.Extensions.Logging;
namespace CryptoExchange.Net.Objects
{
/// <summary>
/// The time synchronization state of an API client
/// </summary>
public class TimeSyncState
{
/// <summary>
/// Name of the API
/// </summary>
public string ApiName { get; set; }
/// <summary>
/// Semaphore to use for checking the time syncing. Should be shared instance among the API client
/// </summary>
public SemaphoreSlim Semaphore { get; }
/// <summary>
/// Last sync time for the API client
/// </summary>
public DateTime LastSyncTime { get; set; }
/// <summary>
/// Time offset for the API client
/// </summary>
public TimeSpan TimeOffset { get; set; }
/// <summary>
/// ctor
/// </summary>
public TimeSyncState(string apiName)
{
ApiName = apiName;
Semaphore = new SemaphoreSlim(1, 1);
}
}
/// <summary>
/// Time synchronization info
/// </summary>
public class TimeSyncInfo
{
/// <summary>
/// Logger
/// </summary>
public ILogger Logger { get; }
/// <summary>
/// Should synchronize time
/// </summary>
public bool SyncTime { get; }
/// <summary>
/// Timestamp recalulcation interval
/// </summary>
public TimeSpan RecalculationInterval { get; }
/// <summary>
/// Time sync state for the API client
/// </summary>
public TimeSyncState TimeSyncState { get; }
/// <summary>
/// ctor
/// </summary>
/// <param name="logger"></param>
/// <param name="recalculationInterval"></param>
/// <param name="syncTime"></param>
/// <param name="syncState"></param>
public TimeSyncInfo(ILogger logger, bool syncTime, TimeSpan recalculationInterval, TimeSyncState syncState)
{
Logger = logger;
SyncTime = syncTime;
RecalculationInterval = recalculationInterval;
TimeSyncState = syncState;
}
/// <summary>
/// Set the time offset
/// </summary>
/// <param name="offset"></param>
public void UpdateTimeOffset(TimeSpan offset)
{
TimeSyncState.LastSyncTime = DateTime.UtcNow;
if (offset.TotalMilliseconds > 0 && offset.TotalMilliseconds < 500)
{
Logger.Log(LogLevel.Information, "{TimeSyncState.ApiName} Time offset within limits, set offset to 0ms", TimeSyncState.ApiName);
TimeSyncState.TimeOffset = TimeSpan.Zero;
}
else
{
Logger.Log(LogLevel.Information, "{TimeSyncState.ApiName} Time offset set to {Offset}ms", TimeSyncState.ApiName, Math.Round(offset.TotalMilliseconds));
TimeSyncState.TimeOffset = offset;
}
}
}
}
@@ -3,24 +3,28 @@ using System;
namespace CryptoExchange.Net.OrderBook
{
internal class ProcessQueueItem
internal class OrderBookUpdate
{
public long StartUpdateId { get; set; }
public long EndUpdateId { get; set; }
public DateTime? LocalDataTime { get; set; }
public DateTime? ServerDataTime { get; set; }
public long StartSequenceNumber { get; set; }
public long EndSequenceNumber { get; set; }
public ISymbolOrderBookEntry[] Bids { get; set; } = Array.Empty<ISymbolOrderBookEntry>();
public ISymbolOrderBookEntry[] Asks { get; set; } = Array.Empty<ISymbolOrderBookEntry>();
}
internal class InitialOrderBookItem
internal class OrderBookSnapshot
{
public long StartUpdateId { get; set; }
public long EndUpdateId { get; set; }
public DateTime? LocalDataTime { get; set; }
public DateTime? ServerDataTime { get; set; }
public long? SequenceNumber { get; set; }
public ISymbolOrderBookEntry[] Bids { get; set; } = Array.Empty<ISymbolOrderBookEntry>();
public ISymbolOrderBookEntry[] Asks { get; set; } = Array.Empty<ISymbolOrderBookEntry>();
}
internal class ChecksumItem
internal class OrderBookChecksum
{
public long? SequenceNumber { get; set; }
public int Checksum { get; set; }
}
}
+262 -64
View File
@@ -1,6 +1,7 @@
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Diagnostics;
using System.Globalization;
using System.Linq;
using System.Text;
@@ -38,6 +39,7 @@ namespace CryptoExchange.Net.OrderBook
private readonly AsyncResetEvent _queueEvent;
private readonly ConcurrentQueue<object> _processQueue;
private bool _validateChecksum;
private bool _firstUpdateAfterSnapshotDone;
private class EmptySymbolOrderBookEntry : ISymbolOrderBookEntry
{
@@ -49,6 +51,13 @@ namespace CryptoExchange.Net.OrderBook
private static readonly ISymbolOrderBookEntry _emptySymbolOrderBookEntry = new EmptySymbolOrderBookEntry();
private enum SequenceNumberResult
{
Skip,
Ok,
OutOfSync
}
/// <summary>
/// A buffer to store messages received before the initial book snapshot is processed. These messages
/// will be processed after the book snapshot is set. Any messages in this buffer with sequence numbers lower
@@ -76,7 +85,12 @@ namespace CryptoExchange.Net.OrderBook
/// the book will resynchronize as it is deemed out of sync
/// </summary>
protected bool _sequencesAreConsecutive;
/// <summary>
/// Whether the first update message after a snapshot may have overlapping sequence numbers instead of the snapshot sequence number + 1
/// </summary>
protected bool _skipSequenceCheckFirstUpdateAfterSnapshotSet;
/// <summary>
/// Whether levels should be strictly enforced. For example, when an order book has 25 levels and a new update comes in which pushes
/// the current level 25 ask out of the top 25, should the level 26 entry be removed from the book or does the server handle this
@@ -133,6 +147,15 @@ namespace CryptoExchange.Net.OrderBook
/// <inheritdoc/>
public DateTime UpdateTime { get; private set; }
/// <inheritdoc/>
public DateTime? UpdateServerTime { get; private set; }
/// <inheritdoc/>
public DateTime? UpdateLocalTime { get; set; }
/// <inheritdoc/>
public TimeSpan? DataAge => DateTime.UtcNow - UpdateLocalTime;
/// <inheritdoc/>
public int AskCount { get; private set; }
@@ -257,6 +280,7 @@ namespace CryptoExchange.Net.OrderBook
_processBuffer.Clear();
_bookSet = false;
_firstUpdateAfterSnapshotDone = false;
Status = OrderBookStatus.Connecting;
_processTask = Task.Factory.StartNew(ProcessQueue, TaskCreationOptions.LongRunning);
@@ -406,45 +430,91 @@ namespace CryptoExchange.Net.OrderBook
/// Implementation for validating a checksum value with the current order book. If checksum validation fails (returns false)
/// the order book will be resynchronized
/// </summary>
/// <param name="checksum"></param>
/// <returns></returns>
protected virtual bool DoChecksum(int checksum) => true;
/// <summary>
/// Set the initial data for the order book. Typically the snapshot which was requested from the Rest API, or the first snapshot
/// received from a socket subscription
/// Set snapshot data for the order book. Typically the snapshot which was requested from the Rest API, or the first snapshot
/// received from a socket subscription. Will clear any previous data.
/// </summary>
/// <param name="orderBookSequenceNumber">The last update sequence number until which the snapshot is in sync</param>
/// <param name="askList">List of asks</param>
/// <param name="bidList">List of bids</param>
protected void SetInitialOrderBook(long orderBookSequenceNumber, ISymbolOrderBookEntry[] bidList, ISymbolOrderBookEntry[] askList)
/// <param name="serverDataTime">Server data timestamp</param>
/// <param name="localDataTime">local data timestamp</param>
protected void SetSnapshot(
long? orderBookSequenceNumber,
ISymbolOrderBookEntry[] bidList,
ISymbolOrderBookEntry[] askList,
DateTime? serverDataTime = null,
DateTime? localDataTime = null)
{
_processQueue.Enqueue(new InitialOrderBookItem { StartUpdateId = orderBookSequenceNumber, EndUpdateId = orderBookSequenceNumber, Asks = askList, Bids = bidList });
_processQueue.Enqueue(
new OrderBookSnapshot
{
LocalDataTime = localDataTime,
ServerDataTime = serverDataTime,
SequenceNumber = orderBookSequenceNumber,
Asks = askList,
Bids = bidList
});
_queueEvent.Set();
}
/// <summary>
/// Add an update to the process queue. Updates the book by providing changed bids and asks, along with an update number which should be higher than the previous update numbers
/// </summary>
/// <param name="updateId">The sequence number</param>
/// <param name="sequenceNumber">The sequence number</param>
/// <param name="bids">List of updated/new bids</param>
/// <param name="asks">List of updated/new asks</param>
protected void UpdateOrderBook(long updateId, ISymbolOrderBookEntry[] bids, ISymbolOrderBookEntry[] asks)
/// <param name="serverDataTime">Server data timestamp</param>
/// <param name="localDataTime">local data timestamp</param>
protected void UpdateOrderBook(
long sequenceNumber,
ISymbolOrderBookEntry[] bids,
ISymbolOrderBookEntry[] asks,
DateTime? serverDataTime = null,
DateTime? localDataTime = null)
{
_processQueue.Enqueue(new ProcessQueueItem { StartUpdateId = updateId, EndUpdateId = updateId, Asks = asks, Bids = bids });
_processQueue.Enqueue(
new OrderBookUpdate
{
LocalDataTime = localDataTime,
ServerDataTime = serverDataTime,
StartSequenceNumber = sequenceNumber,
EndSequenceNumber = sequenceNumber,
Asks = asks,
Bids = bids
});
_queueEvent.Set();
}
/// <summary>
/// Add an update to the process queue. Updates the book by providing changed bids and asks, along with the first and last sequence number in the update
/// </summary>
/// <param name="firstUpdateId">The sequence number of the first update</param>
/// <param name="lastUpdateId">The sequence number of the last update</param>
/// <param name="firstSequenceNumber">The sequence number of the first update</param>
/// <param name="lastSequenceNumber">The sequence number of the last update</param>
/// <param name="bids">List of updated/new bids</param>
/// <param name="asks">List of updated/new asks</param>
protected void UpdateOrderBook(long firstUpdateId, long lastUpdateId, ISymbolOrderBookEntry[] bids, ISymbolOrderBookEntry[] asks)
/// <param name="serverDataTime">Server data timestamp</param>
/// <param name="localDataTime">local data timestamp</param>
protected void UpdateOrderBook(
long firstSequenceNumber,
long lastSequenceNumber,
ISymbolOrderBookEntry[] bids,
ISymbolOrderBookEntry[] asks,
DateTime? serverDataTime = null,
DateTime? localDataTime = null)
{
_processQueue.Enqueue(new ProcessQueueItem { StartUpdateId = firstUpdateId, EndUpdateId = lastUpdateId, Asks = asks, Bids = bids });
_processQueue.Enqueue(
new OrderBookUpdate
{
LocalDataTime = localDataTime,
ServerDataTime = serverDataTime,
StartSequenceNumber = firstSequenceNumber,
EndSequenceNumber = lastSequenceNumber,
Asks = asks,
Bids = bids
});
_queueEvent.Set();
}
@@ -453,12 +523,27 @@ namespace CryptoExchange.Net.OrderBook
/// </summary>
/// <param name="bids">List of updated/new bids</param>
/// <param name="asks">List of updated/new asks</param>
protected void UpdateOrderBook(ISymbolOrderSequencedBookEntry[] bids, ISymbolOrderSequencedBookEntry[] asks)
/// <param name="serverDataTime">Server data timestamp</param>
/// <param name="localDataTime">local data timestamp</param>
protected void UpdateOrderBook(
ISymbolOrderSequencedBookEntry[] bids,
ISymbolOrderSequencedBookEntry[] asks,
DateTime? serverDataTime = null,
DateTime? localDataTime = null)
{
var highest = Math.Max(bids.Any() ? bids.Max(b => b.Sequence) : 0, asks.Any() ? asks.Max(a => a.Sequence) : 0);
var lowest = Math.Min(bids.Any() ? bids.Min(b => b.Sequence) : long.MaxValue, asks.Any() ? asks.Min(a => a.Sequence) : long.MaxValue);
_processQueue.Enqueue(new ProcessQueueItem { StartUpdateId = lowest, EndUpdateId = highest, Asks = asks, Bids = bids });
_processQueue.Enqueue(
new OrderBookUpdate
{
LocalDataTime = localDataTime,
ServerDataTime = serverDataTime,
StartSequenceNumber = lowest,
EndSequenceNumber = highest,
Asks = asks,
Bids = bids
});
_queueEvent.Set();
}
@@ -466,9 +551,10 @@ namespace CryptoExchange.Net.OrderBook
/// Add a checksum value to the process queue
/// </summary>
/// <param name="checksum">The checksum value</param>
protected void AddChecksum(int checksum)
/// <param name="sequenceNumber">The sequence number of the message if it's a separate message with separate number</param>
protected void AddChecksum(int checksum, long? sequenceNumber = null)
{
_processQueue.Enqueue(new ChecksumItem() { Checksum = checksum });
_processQueue.Enqueue(new OrderBookChecksum() { Checksum = checksum, SequenceNumber = sequenceNumber });
_queueEvent.Set();
}
@@ -481,7 +567,12 @@ namespace CryptoExchange.Net.OrderBook
_logger.OrderBookProcessingBufferedUpdates(Api, Symbol, _processBuffer.Count);
foreach (var bufferEntry in _processBuffer)
ProcessRangeUpdates(bufferEntry.FirstUpdateId, bufferEntry.LastUpdateId, bufferEntry.Bids, bufferEntry.Asks);
{
if (_stopProcessing)
break;
ProcessUpdate(bufferEntry.FirstUpdateId, bufferEntry.LastUpdateId, bufferEntry.Bids, bufferEntry.Asks, true);
}
_processBuffer.Clear();
}
@@ -489,26 +580,10 @@ namespace CryptoExchange.Net.OrderBook
/// <summary>
/// Update order book with an entry
/// </summary>
/// <param name="sequence">Sequence number of the update</param>
/// <param name="type">Type of entry</param>
/// <param name="entry">The entry</param>
protected virtual bool ProcessUpdate(long sequence, OrderBookEntryType type, ISymbolOrderBookEntry entry)
protected virtual bool UpdateValue(OrderBookEntryType type, ISymbolOrderBookEntry entry)
{
if (sequence <= LastSequenceNumber)
{
_logger.OrderBookSkippedMessage(Api, Symbol, sequence, LastSequenceNumber);
return false;
}
if (_sequencesAreConsecutive && sequence > LastSequenceNumber + 1)
{
// Out of sync
_logger.OrderBookOutOfSync(Api, Symbol, LastSequenceNumber + 1, sequence);
_stopProcessing = true;
Resubscribe();
return false;
}
UpdateTime = DateTime.UtcNow;
var listToChange = type == OrderBookEntryType.Ask ? _asks : _bids;
if (entry.Quantity == 0)
@@ -565,6 +640,34 @@ namespace CryptoExchange.Net.OrderBook
return new CallResult<bool>(true);
}
/// <summary>
/// Wait until an update has been buffered
/// </summary>
/// <param name="timeout">Max wait time</param>
/// <param name="ct">Cancellation token</param>
/// <returns></returns>
protected async Task<CallResult<bool>> WaitUntilFirstUpdateBufferedAsync(TimeSpan timeout, CancellationToken ct)
{
var startWait = DateTime.UtcNow;
while (_processBuffer.Count == 0)
{
if (ct.IsCancellationRequested)
return new CallResult<bool>(new CancellationRequestedError());
if (DateTime.UtcNow - startWait > timeout)
return new CallResult<bool>(new ServerError(new ErrorInfo(ErrorType.OrderBookTimeout, "Timeout while waiting for data")));
try
{
await Task.Delay(20, ct).ConfigureAwait(false);
}
catch (OperationCanceledException)
{ }
}
return new CallResult<bool>(true);
}
/// <summary>
/// IDisposable implementation for the order book
/// </summary>
@@ -614,6 +717,12 @@ namespace CryptoExchange.Net.OrderBook
{
var stringBuilder = new StringBuilder();
var book = Book;
stringBuilder.AppendLine($"{Exchange} - {Symbol}");
stringBuilder.AppendLine($"Update time local: {UpdateTime:HH:mm:ss.fff} ({Math.Round((DateTime.UtcNow - UpdateTime).TotalMilliseconds)}ms ago)");
stringBuilder.AppendLine($"Data timestamp server: {UpdateServerTime:HH:mm:ss.fff}");
stringBuilder.AppendLine($"Data timestamp local: {UpdateLocalTime:HH:mm:ss.fff}");
stringBuilder.AppendLine($"Data age: {DataAge?.TotalMilliseconds}ms");
stringBuilder.AppendLine();
stringBuilder.AppendLine($" Ask quantity Ask price | Bid price Bid quantity");
for(var i = 0; i < numberOfEntries; i++)
{
@@ -625,6 +734,22 @@ namespace CryptoExchange.Net.OrderBook
return stringBuilder.ToString();
}
/// <inheritdoc />
public Task OutputToConsoleAsync(int numberOfEntries, TimeSpan refreshInterval, CancellationToken ct = default)
{
return Task.Run(async () =>
{
var referenceTime = DateTime.UtcNow;
while (!ct.IsCancellationRequested)
{
Console.Clear();
Console.WriteLine(ToString(numberOfEntries));
var delay = Math.Max(1, (DateTime.UtcNow - referenceTime).TotalMilliseconds % refreshInterval.TotalMilliseconds);
try { await Task.Delay(refreshInterval.Add(TimeSpan.FromMilliseconds(-delay)), ct).ConfigureAwait(false); } catch { }
}
});
}
private void CheckBestOffersChanged(ISymbolOrderBookEntry prevBestBid, ISymbolOrderBookEntry prevBestAsk)
{
var (bestBid, bestAsk) = BestOffers;
@@ -641,8 +766,11 @@ namespace CryptoExchange.Net.OrderBook
// Clear queue
while (_processQueue.TryDequeue(out _)) { }
LastSequenceNumber = 0;
_processBuffer.Clear();
_bookSet = false;
_firstUpdateAfterSnapshotDone = false;
DoReset();
}
@@ -680,17 +808,17 @@ namespace CryptoExchange.Net.OrderBook
continue;
}
if (item is InitialOrderBookItem iobi)
ProcessInitialOrderBookItem(iobi);
if (item is ProcessQueueItem pqi)
ProcessQueueItem(pqi);
else if (item is ChecksumItem ci)
ProcessChecksum(ci);
if (item is OrderBookSnapshot snapshot)
ProcessOrderBookSnapshot(snapshot);
if (item is OrderBookUpdate update)
ProcessQueueItem(update);
else if (item is OrderBookChecksum checksum)
ProcessChecksum(checksum);
}
}
}
private void ProcessInitialOrderBookItem(InitialOrderBookItem item)
private void ProcessOrderBookSnapshot(OrderBookSnapshot item)
{
lock (_bookLock)
{
@@ -702,20 +830,25 @@ namespace CryptoExchange.Net.OrderBook
foreach (var bid in item.Bids)
_bids.Add(bid.Price, bid);
LastSequenceNumber = item.EndUpdateId;
if (item.SequenceNumber != null)
LastSequenceNumber = item.SequenceNumber.Value;
AskCount = _asks.Count;
BidCount = _bids.Count;
UpdateTime = DateTime.UtcNow;
_logger.OrderBookDataSet(Api, Symbol, BidCount, AskCount, item.EndUpdateId);
UpdateServerTime = item.ServerDataTime;
UpdateLocalTime = item.LocalDataTime;
_logger.OrderBookDataSet(Api, Symbol, BidCount, AskCount, item.SequenceNumber);
CheckProcessBuffer();
OnOrderBookUpdate?.Invoke((item.Bids.ToArray(), item.Asks.ToArray()));
OnBestOffersChanged?.Invoke((BestBid, BestAsk));
}
}
private void ProcessQueueItem(ProcessQueueItem item)
private void ProcessQueueItem(OrderBookUpdate item)
{
lock (_bookLock)
{
@@ -725,19 +858,19 @@ namespace CryptoExchange.Net.OrderBook
{
Asks = item.Asks,
Bids = item.Bids,
FirstUpdateId = item.StartUpdateId,
LastUpdateId = item.EndUpdateId,
FirstUpdateId = item.StartSequenceNumber,
LastUpdateId = item.EndSequenceNumber,
});
if (_logger.IsEnabled(LogLevel.Trace))
_logger.OrderBookUpdateBuffered(Api, Symbol, item.StartUpdateId, item.EndUpdateId, item.Asks.Length, item.Bids.Length);
_logger.OrderBookUpdateBuffered(Api, Symbol, item.StartSequenceNumber, item.EndSequenceNumber, item.Asks.Length, item.Bids.Length);
}
else
{
CheckProcessBuffer();
var (prevBestBid, prevBestAsk) = BestOffers;
ProcessRangeUpdates(item.StartUpdateId, item.EndUpdateId, item.Bids, item.Asks);
ProcessUpdate(item.StartSequenceNumber, item.EndSequenceNumber, item.Bids, item.Asks, false);
if (_asks.Count == 0 || _bids.Count == 0)
return;
@@ -750,13 +883,16 @@ namespace CryptoExchange.Net.OrderBook
return;
}
UpdateServerTime = item.ServerDataTime;
UpdateLocalTime = item.LocalDataTime;
OnOrderBookUpdate?.Invoke((item.Bids.ToArray(), item.Asks.ToArray()));
CheckBestOffersChanged(prevBestBid, prevBestAsk);
}
}
}
private void ProcessChecksum(ChecksumItem ci)
private void ProcessChecksum(OrderBookChecksum ci)
{
lock (_bookLock)
{
@@ -776,6 +912,9 @@ namespace CryptoExchange.Net.OrderBook
throw;
}
if (ci.SequenceNumber != null)
LastSequenceNumber = ci.SequenceNumber.Value;
if (!checksumResult)
{
_logger.OrderBookOutOfSyncChecksum(Api, Symbol);
@@ -813,40 +952,99 @@ namespace CryptoExchange.Net.OrderBook
});
}
private void ProcessRangeUpdates(long firstUpdateId, long lastUpdateId, IEnumerable<ISymbolOrderBookEntry> bids, IEnumerable<ISymbolOrderBookEntry> asks)
private void ProcessUpdate(
long updateSequenceNumberStart,
long updateSequenceNumberEnd,
IEnumerable<ISymbolOrderBookEntry> bids,
IEnumerable<ISymbolOrderBookEntry> asks,
bool fromBuffer)
{
if (lastUpdateId <= LastSequenceNumber)
var sequenceResult = fromBuffer ? ValidateBufferSequenceNumber(updateSequenceNumberStart, updateSequenceNumberEnd) : ValidateLiveSequenceNumber(updateSequenceNumberStart);
if (sequenceResult == SequenceNumberResult.Skip)
{
_logger.OrderBookUpdateSkipped(Api, Symbol, lastUpdateId, LastSequenceNumber);
if (updateSequenceNumberStart != updateSequenceNumberEnd)
_logger.OrderBookUpdateSkipped(Api, Symbol, updateSequenceNumberStart, updateSequenceNumberEnd, LastSequenceNumber);
else
_logger.OrderBookUpdateSkipped(Api, Symbol, updateSequenceNumberStart, LastSequenceNumber);
return;
}
if (sequenceResult == SequenceNumberResult.OutOfSync)
{
_logger.OrderBookOutOfSync(Api, Symbol, LastSequenceNumber + 1, updateSequenceNumberStart);
_stopProcessing = true;
Resubscribe();
return;
}
foreach (var entry in bids)
ProcessUpdate(LastSequenceNumber + 1, OrderBookEntryType.Bid, entry);
UpdateValue(OrderBookEntryType.Bid, entry);
foreach (var entry in asks)
ProcessUpdate(LastSequenceNumber + 1, OrderBookEntryType.Ask, entry);
UpdateValue(OrderBookEntryType.Ask, entry);
if (Levels.HasValue && _strictLevels)
{
while (this._bids.Count > Levels.Value)
while (_bids.Count > Levels.Value)
{
BidCount--;
this._bids.Remove(this._bids.Last().Key);
_bids.Remove(_bids.Last().Key);
}
while (this._asks.Count > Levels.Value)
while (_asks.Count > Levels.Value)
{
AskCount--;
this._asks.Remove(this._asks.Last().Key);
_asks.Remove(this._asks.Last().Key);
}
}
LastSequenceNumber = lastUpdateId;
_firstUpdateAfterSnapshotDone = true;
LastSequenceNumber = updateSequenceNumberEnd;
if (_logger.IsEnabled(LogLevel.Trace))
_logger.OrderBookProcessedMessage(Api, Symbol, firstUpdateId, lastUpdateId);
}
{
if (updateSequenceNumberStart != updateSequenceNumberEnd)
_logger.OrderBookProcessedMessage(Api, Symbol, updateSequenceNumberStart, updateSequenceNumberEnd);
else
_logger.OrderBookProcessedMessage(Api, Symbol, updateSequenceNumberStart);
}
}
private SequenceNumberResult ValidateBufferSequenceNumber(long startSequenceNumber, long endSequenceNumber)
{
if (endSequenceNumber <= LastSequenceNumber)
// Buffered update is from before the snapshot, ignore
return SequenceNumberResult.Skip;
if (_sequencesAreConsecutive && startSequenceNumber != LastSequenceNumber + 1)
{
if (_firstUpdateAfterSnapshotDone || !_skipSequenceCheckFirstUpdateAfterSnapshotSet)
// Buffered update is not the next sequence number when it was expected to be
return SequenceNumberResult.OutOfSync;
}
// Buffered sequence number is larger than the last sequence number
return SequenceNumberResult.Ok;
}
private SequenceNumberResult ValidateLiveSequenceNumber(long sequenceNumber)
{
if (sequenceNumber < LastSequenceNumber
&& (_firstUpdateAfterSnapshotDone || !_skipSequenceCheckFirstUpdateAfterSnapshotSet))
// Update is somehow from before the current state
return SequenceNumberResult.OutOfSync;
if (_sequencesAreConsecutive
&& LastSequenceNumber != 0
&& sequenceNumber != LastSequenceNumber + 1)
{
if (_firstUpdateAfterSnapshotDone || !_skipSequenceCheckFirstUpdateAfterSnapshotSet)
return SequenceNumberResult.OutOfSync;
}
return SequenceNumberResult.Ok;
}
}
internal class DescComparer<T> : IComparer<T>
@@ -257,6 +257,7 @@ namespace CryptoExchange.Net.Sockets.Default
}
}
private bool _pausedActivity;
#if NET9_0_OR_GREATER
private readonly Lock _listenersLock = new Lock();
@@ -274,6 +275,8 @@ namespace CryptoExchange.Net.Sockets.Default
private ISocketMessageHandler? _byteMessageConverter;
private ISocketMessageHandler? _textMessageConverter;
private long _lastSequenceNumber;
/// <summary>
/// The task that is sending periodic data on the websocket. Can be used for sending Ping messages every x seconds or similar. Not necessary.
/// </summary>
@@ -293,7 +296,7 @@ namespace CryptoExchange.Net.Sockets.Default
/// Cache for deserialization, only caches for a single message
/// </summary>
private readonly Dictionary<Type, object> _deserializationCache = new Dictionary<Type, object>();
/// <summary>
/// New socket connection
/// </summary>
@@ -340,6 +343,7 @@ namespace CryptoExchange.Net.Sockets.Default
{
Status = SocketStatus.Closed;
Authenticated = false;
_lastSequenceNumber = 0;
if (ApiClient._socketConnections.ContainsKey(SocketId))
ApiClient._socketConnections.TryRemove(SocketId, out _);
@@ -371,6 +375,7 @@ namespace CryptoExchange.Net.Sockets.Default
Status = SocketStatus.Reconnecting;
DisconnectTime = DateTime.UtcNow;
Authenticated = false;
_lastSequenceNumber = 0;
lock (_listenersLock)
{
@@ -566,8 +571,12 @@ namespace CryptoExchange.Net.Sockets.Default
if (deserializationType == null)
{
// No handler found for identifier either, can't process
_logger.LogWarning("Failed to determine message type for identifier {Identifier}. Data: {Message}", typeIdentifier, Encoding.UTF8.GetString(data.ToArray()));
if (!ApiClient.HandleUnhandledMessage(this, typeIdentifier, data))
{
// No handler found for identifier either, can't process
_logger.LogWarning("Failed to determine message type for identifier {Identifier}. Data: {Message}", typeIdentifier, Encoding.UTF8.GetString(data.ToArray()));
}
return;
}
@@ -877,8 +886,16 @@ namespace CryptoExchange.Net.Sockets.Default
subscription.CancellationTokenRegistration.Value.Dispose();
bool anyDuplicateSubscription;
lock (_listenersLock)
anyDuplicateSubscription = _listeners.OfType<Subscription>().Any(x => x != subscription && x.MessageMatcher.HandlerLinks.All(l => subscription.MessageMatcher.ContainsCheck(l)));
if (ApiClient.ClientOptions.UseUpdatedDeserialization)
{
lock (_listenersLock)
anyDuplicateSubscription = _listeners.OfType<Subscription>().Any(x => x != subscription && x.MessageRouter.Routes.All(l => subscription.MessageRouter.ContainsCheck(l)));
}
else
{
lock (_listenersLock)
anyDuplicateSubscription = _listeners.OfType<Subscription>().Any(x => x != subscription && x.MessageMatcher.HandlerLinks.All(l => subscription.MessageMatcher.ContainsCheck(l)));
}
bool shouldCloseConnection;
lock (_listenersLock)
@@ -1228,7 +1245,7 @@ namespace CryptoExchange.Net.Sockets.Default
subQuery.OnComplete = () =>
{
subscription.Status = subQuery.Result!.Success ? SubscriptionStatus.Subscribed : SubscriptionStatus.Pending;
subscription.HandleSubQueryResponse(subQuery.Response);
subscription.HandleSubQueryResponse(this, subQuery.Response);
};
taskList.Add(SendAndWaitQueryAsync(subQuery));
@@ -1268,10 +1285,28 @@ namespace CryptoExchange.Net.Sockets.Default
return CallResult.SuccessResult;
var result = await SendAndWaitQueryAsync(subQuery).ConfigureAwait(false);
subscription.HandleSubQueryResponse(subQuery.Response!);
subscription.HandleSubQueryResponse(this, subQuery.Response!);
return result;
}
/// <summary>
/// Update the sequence number for this connection
/// </summary>
public void UpdateSequenceNumber(long sequenceNumber)
{
if (ApiClient.EnforceSequenceNumbers
&& _lastSequenceNumber != 0 // Initial value is 0
&& _lastSequenceNumber != sequenceNumber // When there are multiple listeners for the same message it's possible this gets recorded multiple times, shouldn't be an issue
&& _lastSequenceNumber + 1 != sequenceNumber) // Expected value
{
// Not sequential
_logger.LogWarning("[Sckt {SocketId}] update not in sequence. Last recorded sequence number: {LastSequence}, update sequence number: {UpdateSequence}. Reconnecting", SocketId, _lastSequenceNumber, sequenceNumber);
_ = TriggerReconnectAsync();
}
_lastSequenceNumber = sequenceNumber;
}
/// <summary>
/// Periodically sends data over a socket connection
/// </summary>
@@ -149,14 +149,12 @@ namespace CryptoExchange.Net.Sockets.Default
/// <summary>
/// Handle a subscription query response
/// </summary>
/// <param name="message"></param>
public virtual void HandleSubQueryResponse(object? message) { }
public virtual void HandleSubQueryResponse(SocketConnection connection, object? message) { }
/// <summary>
/// Handle an unsubscription query response
/// </summary>
/// <param name="message"></param>
public virtual void HandleUnsubQueryResponse(object message) { }
public virtual void HandleUnsubQueryResponse(SocketConnection connection, object message) { }
/// <summary>
/// Create a new unsubscription query
+151
View File
@@ -0,0 +1,151 @@
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
namespace CryptoExchange.Net
{
/// <summary>
/// Manager for timing offsets in APIs
/// </summary>
public static class TimeOffsetManager
{
class SocketTimeOffset
{
private DateTime _lastRollOver = DateTime.UtcNow;
private double? _fallbackLowest;
private double? _currentLowestOffset;
/// <summary>
/// Get the estimated offset, resolves to the lowest offset in time measured in the last two minutes
/// </summary>
public double? Offset
{
get
{
if (_currentLowestOffset == null)
// If there is no current lowest offset return the fallback (which might or might not be null)
return _fallbackLowest;
if (_fallbackLowest == null)
// If there is no fallback return the current lowest offset
return _currentLowestOffset;
// If there is both a fallback and a current offset return the min offset of those
return Math.Min(_currentLowestOffset.Value, _fallbackLowest.Value);
}
}
public void Update(double offsetMs)
{
if (_currentLowestOffset == null || _currentLowestOffset > offsetMs)
{
_currentLowestOffset = offsetMs;
_fallbackLowest = offsetMs;
}
if (DateTime.UtcNow - _lastRollOver > TimeSpan.FromMinutes(1))
{
_fallbackLowest = _currentLowestOffset;
_currentLowestOffset = null;
_lastRollOver = DateTime.UtcNow;
}
}
}
class RestTimeOffset
{
public SemaphoreSlim SemaphoreSlim { get; } = new SemaphoreSlim(1, 1);
public DateTime? LastUpdate { get; set; }
public double? Offset { get; set; }
public void Update(double offsetMs)
{
LastUpdate = DateTime.UtcNow;
Offset = offsetMs;
}
}
private static ConcurrentDictionary<string, SocketTimeOffset> _lastSocketDelays = new ConcurrentDictionary<string, SocketTimeOffset>();
private static ConcurrentDictionary<string, RestTimeOffset> _lastRestDelays = new ConcurrentDictionary<string, RestTimeOffset>();
/// <summary>
/// Update WebSocket API offset
/// </summary>
/// <param name="api">API name</param>
/// <param name="offsetMs">Offset in milliseconds</param>
public static void UpdateSocketOffset(string api, double offsetMs)
{
if (!_lastSocketDelays.TryGetValue(api, out var offsetValues))
{
offsetValues = new SocketTimeOffset();
_lastSocketDelays.TryAdd(api, offsetValues);
}
_lastSocketDelays[api].Update(offsetMs);
}
/// <summary>
/// Update REST API offset
/// </summary>
/// <param name="api">API name</param>
/// <param name="offsetMs">Offset in milliseconds</param>
public static void UpdateRestOffset(string api, double offsetMs)
{
_lastRestDelays[api].Update(offsetMs);
}
/// <summary>
/// Get REST API offset
/// </summary>
/// <param name="api">API name</param>
public static TimeSpan? GetRestOffset(string api) => _lastRestDelays.TryGetValue(api, out var val) && val.Offset != null ? TimeSpan.FromMilliseconds(val.Offset.Value) : null;
/// <summary>
/// Get REST API last update time
/// </summary>
/// <param name="api">API name</param>
public static DateTime? GetRestLastUpdateTime(string api) => _lastRestDelays.TryGetValue(api, out var val) && val.LastUpdate != null ? val.LastUpdate.Value : null;
/// <summary>
/// Register a REST API client to be tracked
/// </summary>
/// <param name="api"></param>
internal static void RegisterRestApi(string api)
{
if (!_lastRestDelays.ContainsKey(api))
_lastRestDelays.TryAdd(api, new RestTimeOffset());
}
/// <summary>
/// Enter exclusive access for the API to update the time offset
/// </summary>
/// <param name="api"></param>
/// <returns></returns>
public static async ValueTask EnterAsync(string api)
{
await _lastRestDelays[api].SemaphoreSlim.WaitAsync().ConfigureAwait(false);
}
/// <summary>
/// Release exclusive access for the API
/// </summary>
/// <param name="api"></param>
public static void Release(string api) => _lastRestDelays[api].SemaphoreSlim.Release();
/// <summary>
/// Get WebSocket API offset
/// </summary>
/// <param name="api">API name</param>
public static TimeSpan? GetSocketOffset(string api) => _lastSocketDelays.TryGetValue(api, out var val) && val.Offset != null ? TimeSpan.FromMilliseconds(val.Offset.Value) : null;
/// <summary>
/// Reset the WebSocket API update timestamp to trigger a new time offset calculation
/// </summary>
/// <param name="api">API name</param>
public static void ResetRestUpdateTime(string api)
{
_lastRestDelays[api].LastUpdate = null;
}
}
}
+26 -26
View File
@@ -5,32 +5,32 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="11.10.0" />
<PackageReference Include="Bitfinex.Net" Version="9.10.0" />
<PackageReference Include="BitMart.Net" Version="2.11.0" />
<PackageReference Include="BloFin.Net" Version="1.3.0" />
<PackageReference Include="Bybit.Net" Version="5.12.0" />
<PackageReference Include="CoinEx.Net" Version="9.10.0" />
<PackageReference Include="CoinW.Net" Version="1.7.0" />
<PackageReference Include="CryptoCom.Net" Version="2.11.0" />
<PackageReference Include="DeepCoin.Net" Version="2.10.0" />
<PackageReference Include="GateIo.Net" Version="2.12.0" />
<PackageReference Include="HyperLiquid.Net" Version="2.16.0" />
<PackageReference Include="JK.BingX.Net" Version="2.10.0" />
<PackageReference Include="JK.Bitget.Net" Version="2.11.0" />
<PackageReference Include="JK.Mexc.Net" Version="3.11.0" />
<PackageReference Include="JK.OKX.Net" Version="3.10.0" />
<PackageReference Include="Jkorf.Aster.Net" Version="1.2.0" />
<PackageReference Include="JKorf.BitMEX.Net" Version="2.10.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="2.10.0" />
<PackageReference Include="JKorf.HTX.Net" Version="7.10.0" />
<PackageReference Include="JKorf.Upbit.Net" Version="1.1.0" />
<PackageReference Include="KrakenExchange.Net" Version="6.10.0" />
<PackageReference Include="Kucoin.Net" Version="7.10.0" />
<PackageReference Include="Serilog.AspNetCore" Version="9.0.0" />
<PackageReference Include="Toobit.Net" Version="1.9.0" />
<PackageReference Include="WhiteBit.Net" Version="2.11.0" />
<PackageReference Include="XT.Net" Version="2.10.0" />
<PackageReference Include="Binance.Net" Version="12.1.0" />
<PackageReference Include="Bitfinex.Net" Version="10.2.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" />
<PackageReference Include="BloFin.Net" Version="2.1.1" />
<PackageReference Include="Bybit.Net" Version="6.1.0" />
<PackageReference Include="CoinEx.Net" Version="10.1.0" />
<PackageReference Include="CoinW.Net" Version="2.1.1" />
<PackageReference Include="CryptoCom.Net" Version="3.1.0" />
<PackageReference Include="DeepCoin.Net" Version="3.1.0" />
<PackageReference Include="GateIo.Net" Version="3.1.0" />
<PackageReference Include="HyperLiquid.Net" Version="3.2.0" />
<PackageReference Include="JK.BingX.Net" Version="3.1.0" />
<PackageReference Include="JK.Bitget.Net" Version="3.1.0" />
<PackageReference Include="JK.Mexc.Net" Version="4.1.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" />
<PackageReference Include="Jkorf.Aster.Net" Version="2.1.0" />
<PackageReference Include="JKorf.BitMEX.Net" Version="3.1.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="3.1.0" />
<PackageReference Include="JKorf.HTX.Net" Version="8.1.0" />
<PackageReference Include="JKorf.Upbit.Net" Version="2.1.0" />
<PackageReference Include="KrakenExchange.Net" Version="7.1.0" />
<PackageReference Include="Kucoin.Net" Version="8.1.0" />
<PackageReference Include="Serilog.AspNetCore" Version="10.0.0" />
<PackageReference Include="Toobit.Net" Version="2.1.0" />
<PackageReference Include="WhiteBit.Net" Version="3.1.0" />
<PackageReference Include="XT.Net" Version="3.1.0" />
</ItemGroup>
</Project>
+1 -1
View File
@@ -50,7 +50,7 @@
binanceSocketClient.SpotApi.ExchangeData.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Binance", data.Data.LastPrice)),
bingXSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("BingX", data.Data.LastPrice)),
bitfinexSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("tETHBTC", data => UpdateData("Bitfinex", data.Data.LastPrice)),
bitgetSocketClient.SpotApiV2.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bitget", data.Data.LastPrice)),
bitgetSocketClient.SpotApiV2.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bitget", data.Data.First().LastPrice)),
bitmartSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH_BTC", data => UpdateData("BitMart", data.Data.LastPrice)),
bitmexSocketClient.ExchangeApi.SubscribeToSymbolUpdatesAsync("ETH_XBT", data => UpdateData("BitMEX", data.Data.LastPrice ?? 0)),
// BloFin doesn't support the ETH/BTC pair
@@ -62,6 +62,7 @@
{
<div style="margin-bottom: 20px; flex: 1; min-width: 300px;">
<h4>@book.Key - @book.Value.Symbol</h4>
<div>@(book.Value.DataAge?.TotalMilliseconds)ms since last update</div>
@if (book.Value.AskCount >= 3 && book.Value.BidCount >= 3)
{
for (var i = 0; i < 3; i++)
+14 -14
View File
@@ -6,20 +6,20 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="11.10.0" />
<PackageReference Include="Bitfinex.Net" Version="9.10.0" />
<PackageReference Include="BitMart.Net" Version="2.11.0" />
<PackageReference Include="Bybit.Net" Version="5.12.0" />
<PackageReference Include="CoinEx.Net" Version="9.10.0" />
<PackageReference Include="CryptoCom.Net" Version="2.11.0" />
<PackageReference Include="GateIo.Net" Version="2.12.0" />
<PackageReference Include="JK.Bitget.Net" Version="2.11.0" />
<PackageReference Include="JK.Mexc.Net" Version="3.11.0" />
<PackageReference Include="JK.OKX.Net" Version="3.10.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="2.10.0" />
<PackageReference Include="JKorf.HTX.Net" Version="7.10.0" />
<PackageReference Include="KrakenExchange.Net" Version="6.10.0" />
<PackageReference Include="Kucoin.Net" Version="7.10.0" />
<PackageReference Include="Binance.Net" Version="12.1.0" />
<PackageReference Include="Bitfinex.Net" Version="10.2.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" />
<PackageReference Include="Bybit.Net" Version="6.1.0" />
<PackageReference Include="CoinEx.Net" Version="10.1.0" />
<PackageReference Include="CryptoCom.Net" Version="3.1.0" />
<PackageReference Include="GateIo.Net" Version="3.1.0" />
<PackageReference Include="JK.Bitget.Net" Version="3.1.0" />
<PackageReference Include="JK.Mexc.Net" Version="4.1.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="3.1.0" />
<PackageReference Include="JKorf.HTX.Net" Version="8.1.0" />
<PackageReference Include="KrakenExchange.Net" Version="7.1.0" />
<PackageReference Include="Kucoin.Net" Version="8.1.0" />
</ItemGroup>
</Project>
+3 -3
View File
@@ -8,9 +8,9 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="11.10.0" />
<PackageReference Include="BitMart.Net" Version="2.11.0" />
<PackageReference Include="JK.OKX.Net" Version="3.10.0" />
<PackageReference Include="Binance.Net" Version="12.1.0" />
<PackageReference Include="BitMart.Net" Version="3.1.0" />
<PackageReference Include="JK.OKX.Net" Version="4.1.0" />
</ItemGroup>
</Project>
+50
View File
@@ -66,6 +66,56 @@ 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).
## Release notes
* Version 10.2.4 - 17 Jan 2026
* Added WaitUntilFirstUpdateBufferedAsync method on SymbolOrderBook
* Added some util methods
* Added CommaSplitStringConverter
* Fixed sequence validation bug SymbolOrderBook
* Version 10.2.3 - 14 Jan 2026
* Added HandleUnhandledMessage virtual method to SocketApiClient to allow some processing for messages which couldn't be mapped via the normal way
* Fixed semaphore exception when creating a new REST client while time sync is in progress on another client
* Version 10.2.2 - 13 Jan 2026
* Allow the same websocket connection sequence number to be recorded multiple times
* Version 10.2.1 - 13 Jan 2026
* Removed duplicate logging for rest responses in Trace verbosity
* Fixed parameter URL creation for array values with ArrayParametersSerialization.MultipleValues
* Version 10.2.0 - 12 Jan 2026
* Added EnforceSequenceNumbers property on SocketApiClient to configure whether websocket message contain sequence numbers and if these should be checked to be sequential
* Added fallback to existing websocket connection if no dedicated request connection was found
* Added IntBoolConverter base class for arbitrary int value to bool mapping
* Added SequenceNumber property to DataEvent object
* Added _skipSequenceCheckFirstUpdateAfterSnapshotSet property for SymbolOrderBook implementations
* Updated SymbolOrderBook sequenceNumber validation
* Updated SymbolOrderBook log verbosities
* Renamed SetInitialOrderBook to SetSnapshot in SymbolOrderBook
* Renamed updateId references to sequenceNumber in SymbolOrderBook
* Version 10.1.0 - 07 Jan 2026
* Updated time sync / time offset management for REST API's
* Added time offset tracking for WebSocket API's
* Added GetAuthenticationQuery virtual method on AuthenticationProvider
* Updated AuthenticationProvider GetTimestamp methods to include a one second offset by default
* Added AuthenticationProvider GetTimestamp methods for SocketApiClient instances
* Added ClientName property on BaseApiClient, resolving to the type name
* Added ObjectOrArrayConverter JsonConverterFactory implementation for resolving json data which might be returned as object or array
* Added UpdateServerTime, UpdateLocalTime and DataAge properties to (I)SymbolOrderBook
* Added OutputToConsoleAsync method to (I)SymbolOrderBook
* Updated SymbolOrderBook string representation
* Added DataTimeLocal and DataAge properties to DataEvent object
* Added SocketConnection parameter to subscription HandleSubQueryResponse and HandleUnsubQueryResponse methods
* Added some utility methods
* Version 10.0.2 - 19 Dec 2025
* Fixed duplicate subscription check with updated deserialization
* Added exception handlers for REST response processing
* Version 10.0.1 - 18 Dec 2025
* Fixed query array parameter serialization
* Version 10.0.0 - 16 Dec 2025
* -