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

Compare commits

...

26 Commits

Author SHA1 Message Date
Jkorf 15657ba683 Updated to version 8.1.1 2024-11-01 10:38:30 +01:00
Jkorf 1aed9f0c67 Fixed System.Text.Json ArrayConverter not passing serializer options to nested deserialization, fixed creating new serializer options each time a JsonConverter attribute is encountered 2024-11-01 10:34:07 +01:00
Jkorf 17f1560310 Fixed socket connections trying to authenticated connection when it's marked as dedicated request connection even when no authentication is needed 2024-11-01 09:38:01 +01:00
Jkorf 41de0a3150 Update index.html 2024-10-28 16:14:54 +01:00
Jkorf 3e410be611 Update index.html 2024-10-28 16:11:01 +01:00
Jkorf be75449e4a Updated examples, added trackers example 2024-10-28 15:38:25 +01:00
Jkorf b1b05c8f6b Added catch around HttpClientHandler.AutomaticDecompression setting as it's not support on Blazor WASM 2024-10-28 13:41:58 +01:00
Jkorf a0e588c3de Updated to version 8.1.0 2024-10-28 10:44:45 +01:00
Jan Korf 9e86a08327 Trackers (#218)
Fix for intermittently failing rate limiting test
Added ConnectionId to RequestDefinition to correctly handle connection and path rate limiting configuration
Added ValidateMessage method to websocket Query object to filter messages even though it is matched to the query based on the  ListenIdentifier
Added KlineTracker and TradeTracker implementation
2024-10-28 10:36:19 +01:00
Jkorf ed007b5272 Added overload for Create method in OrderBookFactory using SharedSymbol 2024-10-23 14:04:10 +02:00
Jkorf bdd5526244 Added Side to SharedTrade model 2024-10-23 14:01:17 +02:00
Jkorf ce35e30688 Updated documentation and examples 2024-10-22 16:20:27 +02:00
Jkorf b40f72b1b0 Added Crypto.com reference 2024-10-22 15:37:56 +02:00
Jkorf 31a6cf285b Doc fix 2024-10-22 11:56:59 +02:00
Jkorf 1842f4fda0 Set ApiCredentials in the client specific options to prevent unknown client credentials when using SetApiCredentials method on client 2024-10-22 10:15:17 +02:00
Jkorf 7a58902ab6 Made SharedFuturesTicker last/high/low price properties nullable 2024-10-22 09:54:44 +02:00
Jkorf 3cb91296ca Added DoHandleReset method for websocket subscriptions 2024-10-22 09:47:14 +02:00
Jkorf 130ed40580 Comment 2024-10-21 16:29:56 +02:00
Jkorf 94cb2caf0b Added System.Text.Json ArrayConverter Write implementation 2024-10-15 10:51:19 +02:00
Jkorf 917d060827 Updated to version 8.0.3 2024-10-14 14:13:38 +02:00
JKorf c58bc2be07 Merge branch 'master' of https://github.com/JKorf/CryptoExchange.Net 2024-10-12 13:06:28 +02:00
JKorf ff3356e2b4 Added Authenticated property on base client and shared client 2024-10-12 13:06:22 +02:00
Jkorf 79434c7be5 Implemented GetValues System.Text.Json in message accessor 2024-10-11 16:02:17 +02:00
Jkorf 168dabc11f Added fallback for unparsable value in System.Text.Json NumberStringConverter 2024-10-09 15:42:11 +02:00
Jkorf 71ee263683 Added support for duplicate array indexes in System.Text.Json ArrayConverter 2024-10-09 15:41:43 +02:00
Jkorf 7239b9c289 docs 2024-10-09 10:25:57 +02:00
47 changed files with 3092 additions and 108 deletions
@@ -1,4 +1,4 @@
<Project Sdk="Microsoft.NET.Sdk"> <Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup> <PropertyGroup>
<TargetFramework>net8.0</TargetFramework> <TargetFramework>net8.0</TargetFramework>
@@ -6,10 +6,10 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.9.0"></PackageReference> <PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.11.1"></PackageReference>
<PackageReference Include="Moq" Version="4.20.70" /> <PackageReference Include="Moq" Version="4.20.72" />
<PackageReference Include="NUnit" Version="4.1.0"></PackageReference> <PackageReference Include="NUnit" Version="4.2.2"></PackageReference>
<PackageReference Include="NUnit3TestAdapter" Version="4.5.0"></PackageReference> <PackageReference Include="NUnit3TestAdapter" Version="4.6.0"></PackageReference>
</ItemGroup> </ItemGroup>
<ItemGroup> <ItemGroup>
@@ -302,7 +302,7 @@ namespace CryptoExchange.Net.UnitTests
public async Task ApiKeyRateLimiterBasics(string key1, string key2, string endpoint1, string endpoint2, bool expectLimited) public async Task ApiKeyRateLimiterBasics(string key1, string key2, string endpoint1, string endpoint2, bool expectLimited)
{ {
var rateLimiter = new RateLimitGate("Test"); var rateLimiter = new RateLimitGate("Test");
rateLimiter.AddGuard(new RateLimitGuard(RateLimitGuard.PerApiKey, new AuthenticatedEndpointFilter(true), 1, TimeSpan.FromSeconds(0.1), RateLimitWindowType.Fixed)); rateLimiter.AddGuard(new RateLimitGuard(RateLimitGuard.PerApiKey, new AuthenticatedEndpointFilter(true), 1, TimeSpan.FromSeconds(0.1), RateLimitWindowType.Sliding));
var requestDefinition1 = new RequestDefinition(endpoint1, HttpMethod.Get) { Authenticated = key1 != null }; var requestDefinition1 = new RequestDefinition(endpoint1, HttpMethod.Get) { Authenticated = key1 != null };
var requestDefinition2 = new RequestDefinition(endpoint2, HttpMethod.Get) { Authenticated = key2 != null }; var requestDefinition2 = new RequestDefinition(endpoint2, HttpMethod.Get) { Authenticated = key2 != null };
@@ -5,6 +5,8 @@ using NUnit.Framework;
using System; using System;
using System.Text.Json.Serialization; using System.Text.Json.Serialization;
using NUnit.Framework.Legacy; using NUnit.Framework.Legacy;
using CryptoExchange.Net.Converters;
using CryptoExchange.Net.Testing.Comparers;
namespace CryptoExchange.Net.UnitTests namespace CryptoExchange.Net.UnitTests
{ {
@@ -242,6 +244,44 @@ namespace CryptoExchange.Net.UnitTests
var result = JsonSerializer.Deserialize<STJDecimalObject>("{ \"test\": " + value + "}"); var result = JsonSerializer.Deserialize<STJDecimalObject>("{ \"test\": " + value + "}");
Assert.That(result.Test, Is.EqualTo(expected == -999 ? decimal.MaxValue : expected)); Assert.That(result.Test, Is.EqualTo(expected == -999 ? decimal.MaxValue : expected));
} }
[Test()]
public void TestArrayConverter()
{
var data = new Test()
{
Prop1 = 2,
Prop2 = null,
Prop3 = "123",
Prop3Again = "123",
Prop4 = null,
Prop5 = new Test2
{
Prop21 = 3,
Prop22 = "456"
},
Prop6 = new Test3
{
Prop31 = 4,
Prop32 = "789"
},
Prop7 = TestEnum.Two
};
var serialized = JsonSerializer.Serialize(data);
var deserialized = JsonSerializer.Deserialize<Test>(serialized);
Assert.That(deserialized.Prop1, Is.EqualTo(2));
Assert.That(deserialized.Prop2, Is.Null);
Assert.That(deserialized.Prop3, Is.EqualTo("123"));
Assert.That(deserialized.Prop3Again, Is.EqualTo("123"));
Assert.That(deserialized.Prop4, Is.Null);
Assert.That(deserialized.Prop5.Prop21, Is.EqualTo(3));
Assert.That(deserialized.Prop5.Prop22, Is.EqualTo("456"));
Assert.That(deserialized.Prop6.Prop31, Is.EqualTo(4));
Assert.That(deserialized.Prop6.Prop32, Is.EqualTo("789"));
Assert.That(deserialized.Prop7, Is.EqualTo(TestEnum.Two));
}
} }
public class STJDecimalObject public class STJDecimalObject
@@ -281,4 +321,42 @@ namespace CryptoExchange.Net.UnitTests
[JsonConverter(typeof(BoolConverter))] [JsonConverter(typeof(BoolConverter))]
public bool Value { get; set; } public bool Value { get; set; }
} }
[JsonConverter(typeof(ArrayConverter))]
record Test
{
[ArrayProperty(0)]
public int Prop1 { get; set; }
[ArrayProperty(1)]
public int? Prop2 { get; set; }
[ArrayProperty(2)]
public string Prop3 { get; set; }
[ArrayProperty(2)]
public string Prop3Again { get; set; }
[ArrayProperty(3)]
public string Prop4 { get; set; }
[ArrayProperty(4)]
public Test2 Prop5 { get; set; }
[ArrayProperty(5)]
public Test3 Prop6 { get; set; }
[ArrayProperty(6), JsonConverter(typeof(EnumConverter))]
public TestEnum? Prop7 { get; set; }
}
[JsonConverter(typeof(ArrayConverter))]
record Test2
{
[ArrayProperty(0)]
public int Prop21 { get; set; }
[ArrayProperty(1)]
public string Prop22 { get; set; }
}
record Test3
{
[JsonPropertyName("prop31")]
public int Prop31 { get; set; }
[JsonPropertyName("prop32")]
public string Prop32 { get; set; }
}
} }
@@ -38,6 +38,11 @@ namespace CryptoExchange.Net.Clients
/// </summary> /// </summary>
public bool OutputOriginalData { get; } public bool OutputOriginalData { get; }
/// <summary>
/// Whether or not API credentials have been configured for this client. Does not check the credentials are actually valid.
/// </summary>
public bool Authenticated => ApiOptions.ApiCredentials != null || ClientOptions.ApiCredentials != null;
/// <summary> /// <summary>
/// Api options /// Api options
/// </summary> /// </summary>
@@ -83,6 +88,7 @@ namespace CryptoExchange.Net.Clients
/// <inheritdoc /> /// <inheritdoc />
public void SetApiCredentials<T>(T credentials) where T : ApiCredentials public void SetApiCredentials<T>(T credentials) where T : ApiCredentials
{ {
ApiOptions.ApiCredentials = credentials;
if (credentials != null) if (credentials != null)
AuthenticationProvider = CreateAuthenticationProvider(credentials.Copy()); AuthenticationProvider = CreateAuthenticationProvider(credentials.Copy());
} }
+14 -4
View File
@@ -490,11 +490,14 @@ namespace CryptoExchange.Net.Clients
SocketConnection connection; SocketConnection connection;
if (!dedicatedRequestConnection) if (!dedicatedRequestConnection)
{ {
connection = socketQuery.Where(s => !s.Value.DedicatedRequestConnection).OrderBy(s => s.Value.UserSubscriptionCount).FirstOrDefault().Value; connection = socketQuery.Where(s => !s.Value.DedicatedRequestConnection.IsDedicatedRequestConnection).OrderBy(s => s.Value.UserSubscriptionCount).FirstOrDefault().Value;
} }
else else
{ {
connection = socketQuery.Where(s => s.Value.DedicatedRequestConnection).FirstOrDefault().Value; connection = socketQuery.Where(s => s.Value.DedicatedRequestConnection.IsDedicatedRequestConnection).FirstOrDefault().Value;
if (connection != null && !connection.DedicatedRequestConnection.Authenticated)
// Mark dedicated request connection as authenticated if the request is authenticated
connection.DedicatedRequestConnection.Authenticated = authenticated;
} }
if (connection != null) if (connection != null)
@@ -519,7 +522,14 @@ namespace CryptoExchange.Net.Clients
var socketConnection = new SocketConnection(_logger, this, socket, address); var socketConnection = new SocketConnection(_logger, this, socket, address);
socketConnection.UnhandledMessage += HandleUnhandledMessage; socketConnection.UnhandledMessage += HandleUnhandledMessage;
socketConnection.ConnectRateLimitedAsync += HandleConnectRateLimitedAsync; socketConnection.ConnectRateLimitedAsync += HandleConnectRateLimitedAsync;
socketConnection.DedicatedRequestConnection = dedicatedRequestConnection; if (dedicatedRequestConnection)
{
socketConnection.DedicatedRequestConnection = new DedicatedConnectionState
{
IsDedicatedRequestConnection = dedicatedRequestConnection,
Authenticated = authenticated
};
}
foreach (var ptg in PeriodicTaskRegistrations) foreach (var ptg in PeriodicTaskRegistrations)
socketConnection.QueryPeriodic(ptg.Identifier, ptg.Interval, ptg.QueryDelegate, ptg.Callback); socketConnection.QueryPeriodic(ptg.Identifier, ptg.Interval, ptg.QueryDelegate, ptg.Callback);
@@ -652,7 +662,7 @@ namespace CryptoExchange.Net.Clients
var tasks = new List<Task>(); var tasks = new List<Task>();
{ {
var socketList = socketConnections.Values; var socketList = socketConnections.Values;
foreach (var connection in socketList.Where(s => !s.DedicatedRequestConnection)) foreach (var connection in socketList.Where(s => !s.DedicatedRequestConnection.IsDedicatedRequestConnection))
tasks.Add(connection.CloseAsync()); tasks.Add(connection.CloseAsync());
} }
@@ -38,12 +38,67 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
private class ArrayConverterInner<T> : JsonConverter<T> private class ArrayConverterInner<T> : JsonConverter<T>
{ {
private static readonly ConcurrentDictionary<Type, List<ArrayPropertyInfo>> _typeAttributesCache = new ConcurrentDictionary<Type, List<ArrayPropertyInfo>>(); private static readonly ConcurrentDictionary<Type, List<ArrayPropertyInfo>> _typeAttributesCache = new ConcurrentDictionary<Type, List<ArrayPropertyInfo>>();
private static readonly ConcurrentDictionary<Type, JsonSerializerOptions> _converterOptionsCache = new ConcurrentDictionary<Type, JsonSerializerOptions>();
public override void Write(Utf8JsonWriter writer, T value, JsonSerializerOptions options) public override void Write(Utf8JsonWriter writer, T value, JsonSerializerOptions options)
{ {
// TODO if (value == null)
throw new NotImplementedException(); {
writer.WriteNullValue();
return;
}
writer.WriteStartArray();
var valueType = value.GetType();
if (!_typeAttributesCache.TryGetValue(valueType, out var typeAttributes))
typeAttributes = CacheTypeAttributes(valueType);
var ordered = typeAttributes.Where(x => x.ArrayProperty != null).OrderBy(p => p.ArrayProperty.Index);
var last = -1;
foreach (var prop in ordered)
{
if (prop.ArrayProperty.Index == last)
continue;
while (prop.ArrayProperty.Index != last + 1)
{
writer.WriteNullValue();
last += 1;
}
last = prop.ArrayProperty.Index;
var objValue = prop.PropertyInfo.GetValue(value);
if (objValue == null)
{
writer.WriteNullValue();
continue;
}
JsonSerializerOptions? typeOptions = null;
if (prop.JsonConverterType != null)
{
var converter = (JsonConverter)Activator.CreateInstance(prop.JsonConverterType);
typeOptions = new JsonSerializerOptions();
typeOptions.Converters.Clear();
typeOptions.Converters.Add(converter);
}
if (prop.JsonConverterType == null && IsSimple(prop.PropertyInfo.PropertyType))
{
if (prop.PropertyInfo.PropertyType == typeof(string))
writer.WriteStringValue(Convert.ToString(objValue, CultureInfo.InvariantCulture));
else
writer.WriteRawValue(Convert.ToString(objValue, CultureInfo.InvariantCulture));
}
else
{
JsonSerializer.Serialize(writer, objValue, typeOptions ?? options);
}
}
writer.WriteEndArray();
} }
/// <inheritdoc /> /// <inheritdoc />
@@ -53,7 +108,20 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
return default; return default;
var result = Activator.CreateInstance(typeToConvert); var result = Activator.CreateInstance(typeToConvert);
return (T)ParseObject(ref reader, result, typeToConvert); return (T)ParseObject(ref reader, result, typeToConvert, options);
}
private static bool IsSimple(Type type)
{
if (type.IsGenericType && type.GetGenericTypeDefinition() == typeof(Nullable<>))
{
// nullable type, check if the nested type is simple.
return IsSimple(type.GetGenericArguments()[0]);
}
return type.IsPrimitive
|| type.IsEnum
|| type == typeof(string)
|| type == typeof(decimal);
} }
private static List<ArrayPropertyInfo> CacheTypeAttributes(Type type) private static List<ArrayPropertyInfo> CacheTypeAttributes(Type type)
@@ -71,7 +139,7 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
ArrayProperty = att, ArrayProperty = att,
PropertyInfo = property, PropertyInfo = property,
DefaultDeserialization = property.GetCustomAttribute<JsonConversionAttribute>() != null, DefaultDeserialization = property.GetCustomAttribute<JsonConversionAttribute>() != null,
JsonConverterType = property.GetCustomAttribute<JsonConverterAttribute>()?.ConverterType, JsonConverterType = property.GetCustomAttribute<JsonConverterAttribute>()?.ConverterType ?? property.PropertyType.GetCustomAttribute<JsonConverterAttribute>()?.ConverterType,
TargetType = Nullable.GetUnderlyingType(property.PropertyType) ?? property.PropertyType TargetType = Nullable.GetUnderlyingType(property.PropertyType) ?? property.PropertyType
}); });
} }
@@ -80,7 +148,7 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
return attributes; return attributes;
} }
private static object ParseObject(ref Utf8JsonReader reader, object result, Type objectType) private static object ParseObject(ref Utf8JsonReader reader, object result, Type objectType, JsonSerializerOptions options)
{ {
if (reader.TokenType != JsonTokenType.StartArray) if (reader.TokenType != JsonTokenType.StartArray)
throw new Exception("Not an array"); throw new Exception("Not an array");
@@ -94,42 +162,58 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
if (reader.TokenType == JsonTokenType.EndArray) if (reader.TokenType == JsonTokenType.EndArray)
break; break;
var attribute = attributes.SingleOrDefault(a => a.ArrayProperty.Index == index); var indexAttributes = attributes.Where(a => a.ArrayProperty.Index == index);
if (attribute == null) if (!indexAttributes.Any())
{ {
index++; index++;
continue; continue;
} }
var targetType = attribute.TargetType; foreach (var attribute in indexAttributes)
object? value = null;
if (attribute.JsonConverterType != null)
{ {
// Has JsonConverter attribute var targetType = attribute.TargetType;
var options = new JsonSerializerOptions(); object? value = null;
options.Converters.Add((JsonConverter)Activator.CreateInstance(attribute.JsonConverterType)); if (attribute.JsonConverterType != null)
value = JsonDocument.ParseValue(ref reader).Deserialize(targetType, options);
}
else if (attribute.DefaultDeserialization)
{
// Use default deserialization
value = JsonDocument.ParseValue(ref reader).Deserialize(targetType);
}
else
{
value = reader.TokenType switch
{ {
JsonTokenType.Null => null, if (!_converterOptionsCache.TryGetValue(attribute.JsonConverterType, out var newOptions))
JsonTokenType.False => false, {
JsonTokenType.True => true, var converter = (JsonConverter)Activator.CreateInstance(attribute.JsonConverterType);
JsonTokenType.String => reader.GetString(), newOptions = new JsonSerializerOptions
JsonTokenType.Number => reader.GetDecimal(), {
_ => throw new NotImplementedException($"Array deserialization of type {reader.TokenType} not supported"), NumberHandling = SerializerOptions.WithConverters.NumberHandling,
}; PropertyNameCaseInsensitive = SerializerOptions.WithConverters.PropertyNameCaseInsensitive,
Converters = { converter },
};
_converterOptionsCache.TryAdd(attribute.JsonConverterType, newOptions);
}
value = JsonDocument.ParseValue(ref reader).Deserialize(targetType, newOptions);
}
else if (attribute.DefaultDeserialization)
{
// Use default deserialization
value = JsonDocument.ParseValue(ref reader).Deserialize(targetType, options);
}
else
{
value = reader.TokenType switch
{
JsonTokenType.Null => null,
JsonTokenType.False => false,
JsonTokenType.True => true,
JsonTokenType.String => reader.GetString(),
JsonTokenType.Number => reader.GetDecimal(),
JsonTokenType.StartObject => JsonSerializer.Deserialize(ref reader, attribute.TargetType, options),
_ => throw new NotImplementedException($"Array deserialization of type {reader.TokenType} not supported"),
};
}
if (targetType.IsAssignableFrom(value?.GetType()))
attribute.PropertyInfo.SetValue(result, value == null ? null : value);
else
attribute.PropertyInfo.SetValue(result, value == null ? null : Convert.ChangeType(value, targetType, CultureInfo.InvariantCulture));
} }
attribute.PropertyInfo.SetValue(result, value == null ? null : Convert.ChangeType(value, targetType, CultureInfo.InvariantCulture));
index++; index++;
} }
@@ -23,7 +23,14 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
return reader.GetDecimal().ToString(); return reader.GetDecimal().ToString();
} }
return reader.GetString(); try
{
return reader.GetString();
}
catch (Exception)
{
return null;
}
} }
/// <inheritdoc /> /// <inheritdoc />
@@ -50,6 +50,11 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
var info = $"Deserialize JsonException: {ex.Message}, Path: {ex.Path}, LineNumber: {ex.LineNumber}, LinePosition: {ex.BytePositionInLine}"; var info = $"Deserialize JsonException: {ex.Message}, Path: {ex.Path}, LineNumber: {ex.LineNumber}, LinePosition: {ex.BytePositionInLine}";
return new CallResult<object>(new DeserializeError(info, OriginalDataAvailable ? GetOriginalString() : "[Data only available when OutputOriginal = true in client options]")); return new CallResult<object>(new DeserializeError(info, OriginalDataAvailable ? GetOriginalString() : "[Data only available when OutputOriginal = true in client options]"));
} }
catch (Exception ex)
{
var info = $"Deserialize unknown Exception: {ex.Message}";
return new CallResult<object>(new DeserializeError(info, OriginalDataAvailable ? GetOriginalString() : "[Data only available when OutputOriginal = true in client options]"));
}
} }
/// <inheritdoc /> /// <inheritdoc />
@@ -133,7 +138,20 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
} }
/// <inheritdoc /> /// <inheritdoc />
public List<T?>? GetValues<T>(MessagePath path) => throw new NotImplementedException(); public List<T?>? GetValues<T>(MessagePath path)
{
if (!IsJson)
throw new InvalidOperationException("Can't access json data on non-json message");
var value = GetPathNode(path);
if (value == null)
return default;
if (value.Value.ValueKind != JsonValueKind.Array)
return default;
return value.Value.Deserialize<List<T>>()!;
}
private JsonElement? GetPathNode(MessagePath path) private JsonElement? GetPathNode(MessagePath path)
{ {
+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>8.0.2</PackageVersion> <PackageVersion>8.1.1</PackageVersion>
<AssemblyVersion>8.0.2</AssemblyVersion> <AssemblyVersion>8.1.1</AssemblyVersion>
<FileVersion>8.0.2</FileVersion> <FileVersion>8.1.1</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</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</PackageTags>
<RepositoryType>git</RepositoryType> <RepositoryType>git</RepositoryType>
@@ -1,4 +1,5 @@
using CryptoExchange.Net.Objects.Options; using CryptoExchange.Net.Objects.Options;
using CryptoExchange.Net.SharedApis;
using System; using System;
namespace CryptoExchange.Net.Interfaces namespace CryptoExchange.Net.Interfaces
@@ -23,5 +24,12 @@ namespace CryptoExchange.Net.Interfaces
/// <param name="options">Options for the order book</param> /// <param name="options">Options for the order book</param>
/// <returns></returns> /// <returns></returns>
public ISymbolOrderBook Create(string baseAsset, string quoteAsset, Action<TOptions>? options = null); public ISymbolOrderBook Create(string baseAsset, string quoteAsset, Action<TOptions>? options = null);
/// <summary>
/// Create a new order book by base and quote asset names
/// </summary>
/// <param name="symbol">Symbol</param>
/// <param name="options">Options for the order book</param>
/// <returns></returns>
public ISymbolOrderBook Create(SharedSymbol symbol, Action<TOptions>? options = null);
} }
} }
@@ -62,7 +62,7 @@ namespace CryptoExchange.Net.Interfaces
Task UnsubscribeAsync(UpdateSubscription subscription); Task UnsubscribeAsync(UpdateSubscription subscription);
/// <summary> /// <summary>
/// Prepare connections which can subsequently be used for sending websocket requests. /// Prepare connections which can subsequently be used for sending websocket requests. Note that this is not required. If not prepared it will be initialized at the first websocket request.
/// </summary> /// </summary>
/// <returns></returns> /// <returns></returns>
Task<CallResult> PrepareConnectionsAsync(); Task<CallResult> PrepareConnectionsAsync();
@@ -0,0 +1,292 @@
using System;
using CryptoExchange.Net.Objects;
using Microsoft.Extensions.Logging;
namespace CryptoExchange.Net.Logging.Extensions
{
#pragma warning disable CS1591 // Missing XML comment for publicly visible type or member
public static class TrackerLoggingExtensions
{
private static readonly Action<ILogger, string, SyncStatus, SyncStatus, Exception?> _klineTrackerStatusChanged;
private static readonly Action<ILogger, string, Exception?> _klineTrackerStarting;
private static readonly Action<ILogger, string, string, Exception?> _klineTrackerStartFailed;
private static readonly Action<ILogger, string, Exception?> _klineTrackerStarted;
private static readonly Action<ILogger, string, Exception?> _klineTrackerStopping;
private static readonly Action<ILogger, string, Exception?> _klineTrackerStopped;
private static readonly Action<ILogger, string, DateTime, Exception?> _klineTrackerInitialDataSet;
private static readonly Action<ILogger, string, DateTime, Exception?> _klineTrackerKlineUpdated;
private static readonly Action<ILogger, string, DateTime, Exception?> _klineTrackerKlineAdded;
private static readonly Action<ILogger, string, Exception?> _klineTrackerConnectionLost;
private static readonly Action<ILogger, string, Exception?> _klineTrackerConnectionClosed;
private static readonly Action<ILogger, string, Exception?> _klineTrackerConnectionRestored;
private static readonly Action<ILogger, string, SyncStatus, SyncStatus, Exception?> _tradeTrackerStatusChanged;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerStarting;
private static readonly Action<ILogger, string, string, Exception?> _tradeTrackerStartFailed;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerStarted;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerStopping;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerStopped;
private static readonly Action<ILogger, string, int, long, Exception?> _tradeTrackerInitialDataSet;
private static readonly Action<ILogger, string, long, Exception?> _tradeTrackerPreSnapshotSkip;
private static readonly Action<ILogger, string, long, Exception?> _tradeTrackerPreSnapshotApplied;
private static readonly Action<ILogger, string, long, Exception?> _tradeTrackerTradeAdded;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerConnectionLost;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerConnectionClosed;
private static readonly Action<ILogger, string, Exception?> _tradeTrackerConnectionRestored;
static TrackerLoggingExtensions()
{
_klineTrackerStatusChanged = LoggerMessage.Define<string, SyncStatus, SyncStatus>(
LogLevel.Debug,
new EventId(6001, "KlineTrackerStatusChanged"),
"Kline tracker for {Symbol} status changed: {OldStatus} => {NewStatus}");
_klineTrackerStarting = LoggerMessage.Define<string>(
LogLevel.Debug,
new EventId(6002, "KlineTrackerStarting"),
"Kline tracker for {Symbol} starting");
_klineTrackerStartFailed = LoggerMessage.Define<string, string>(
LogLevel.Warning,
new EventId(6003, "KlineTrackerStartFailed"),
"Kline tracker for {Symbol} failed to start: {Error}");
_klineTrackerStarted = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6004, "KlineTrackerStarted"),
"Kline tracker for {Symbol} started");
_klineTrackerStopping = LoggerMessage.Define<string>(
LogLevel.Debug,
new EventId(6005, "KlineTrackerStopping"),
"Kline tracker for {Symbol} stopping");
_klineTrackerStopped = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6006, "KlineTrackerStopped"),
"Kline tracker for {Symbol} stopped");
_klineTrackerInitialDataSet = LoggerMessage.Define<string, DateTime>(
LogLevel.Debug,
new EventId(6007, "KlineTrackerInitialDataSet"),
"Kline tracker for {Symbol} initial data set, last timestamp: {LastTime}");
_klineTrackerKlineUpdated = LoggerMessage.Define<string, DateTime>(
LogLevel.Trace,
new EventId(6008, "KlineTrackerKlineUpdated"),
"Kline tracker for {Symbol} kline updated for open time: {LastTime}");
_klineTrackerKlineAdded = LoggerMessage.Define<string, DateTime>(
LogLevel.Trace,
new EventId(6009, "KlineTrackerKlineAdded"),
"Kline tracker for {Symbol} new kline for open time: {LastTime}");
_klineTrackerConnectionLost = LoggerMessage.Define<string>(
LogLevel.Warning,
new EventId(6010, "KlineTrackerConnectionLost"),
"Kline tracker for {Symbol} connection lost");
_klineTrackerConnectionClosed = LoggerMessage.Define<string>(
LogLevel.Warning,
new EventId(6011, "KlineTrackerConnectionClosed"),
"Kline tracker for {Symbol} disconnected");
_klineTrackerConnectionRestored = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6012, "KlineTrackerConnectionRestored"),
"Kline tracker for {Symbol} successfully resynchronized");
_tradeTrackerStatusChanged = LoggerMessage.Define<string, SyncStatus, SyncStatus>(
LogLevel.Debug,
new EventId(6013, "KlineTrackerStatusChanged"),
"Trade tracker for {Symbol} status changed: {OldStatus} => {NewStatus}");
_tradeTrackerStarting = LoggerMessage.Define<string>(
LogLevel.Debug,
new EventId(6014, "KlineTrackerStarting"),
"Trade tracker for {Symbol} starting");
_tradeTrackerStartFailed = LoggerMessage.Define<string, string>(
LogLevel.Warning,
new EventId(6015, "KlineTrackerStartFailed"),
"Trade tracker for {Symbol} failed to start: {Error}");
_tradeTrackerStarted = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6016, "KlineTrackerStarted"),
"Trade tracker for {Symbol} started");
_tradeTrackerStopping = LoggerMessage.Define<string>(
LogLevel.Debug,
new EventId(6017, "KlineTrackerStopping"),
"Trade tracker for {Symbol} stopping");
_tradeTrackerStopped = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6018, "KlineTrackerStopped"),
"Trade tracker for {Symbol} stopped");
_tradeTrackerInitialDataSet = LoggerMessage.Define<string, int, long>(
LogLevel.Debug,
new EventId(6019, "TradeTrackerInitialDataSet"),
"Trade tracker for {Symbol} snapshot set, Count: {Count}, Last id: {LastId}");
_tradeTrackerPreSnapshotSkip = LoggerMessage.Define<string, long>(
LogLevel.Trace,
new EventId(6020, "TradeTrackerPreSnapshotSkip"),
"Trade tracker for {Symbol} skipping {Id}, already in snapshot");
_tradeTrackerPreSnapshotApplied = LoggerMessage.Define<string, long>(
LogLevel.Trace,
new EventId(6021, "TradeTrackerPreSnapshotApplied"),
"Trade tracker for {Symbol} adding {Id} from pre-snapshot");
_tradeTrackerTradeAdded = LoggerMessage.Define<string, long>(
LogLevel.Trace,
new EventId(6022, "TradeTrackerTradeAdded"),
"Trade tracker for {Symbol} adding trade {Id}");
_tradeTrackerConnectionLost = LoggerMessage.Define<string>(
LogLevel.Warning,
new EventId(6023, "TradeTrackerConnectionLost"),
"Trade tracker for {Symbol} connection lost");
_tradeTrackerConnectionClosed = LoggerMessage.Define<string>(
LogLevel.Warning,
new EventId(6024, "TradeTrackerConnectionClosed"),
"Trade tracker for {Symbol} disconnected");
_tradeTrackerConnectionRestored = LoggerMessage.Define<string>(
LogLevel.Information,
new EventId(6025, "TradeTrackerConnectionRestored"),
"Trade tracker for {Symbol} successfully resynchronized");
}
public static void KlineTrackerStatusChanged(this ILogger logger, string symbol, SyncStatus oldStatus, SyncStatus newStatus)
{
_klineTrackerStatusChanged(logger, symbol, oldStatus, newStatus, null);
}
public static void KlineTrackerStarting(this ILogger logger, string symbol)
{
_klineTrackerStarting(logger, symbol, null);
}
public static void KlineTrackerStartFailed(this ILogger logger, string symbol, string error)
{
_klineTrackerStartFailed(logger, symbol, error, null);
}
public static void KlineTrackerStarted(this ILogger logger, string symbol)
{
_klineTrackerStarted(logger, symbol, null);
}
public static void KlineTrackerStopping(this ILogger logger, string symbol)
{
_klineTrackerStopping(logger, symbol, null);
}
public static void KlineTrackerStopped(this ILogger logger, string symbol)
{
_klineTrackerStopped(logger, symbol, null);
}
public static void KlineTrackerInitialDataSet(this ILogger logger, string symbol, DateTime lastTime)
{
_klineTrackerInitialDataSet(logger, symbol, lastTime, null);
}
public static void KlineTrackerKlineUpdated(this ILogger logger, string symbol, DateTime lastTime)
{
_klineTrackerKlineUpdated(logger, symbol, lastTime, null);
}
public static void KlineTrackerKlineAdded(this ILogger logger, string symbol, DateTime lastTime)
{
_klineTrackerKlineAdded(logger, symbol, lastTime, null);
}
public static void KlineTrackerConnectionLost(this ILogger logger, string symbol)
{
_klineTrackerConnectionLost(logger, symbol, null);
}
public static void KlineTrackerConnectionClosed(this ILogger logger, string symbol)
{
_klineTrackerConnectionClosed(logger, symbol, null);
}
public static void KlineTrackerConnectionRestored(this ILogger logger, string symbol)
{
_klineTrackerConnectionRestored(logger, symbol, null);
}
public static void TradeTrackerStatusChanged(this ILogger logger, string symbol, SyncStatus oldStatus, SyncStatus newStatus)
{
_tradeTrackerStatusChanged(logger, symbol, oldStatus, newStatus, null);
}
public static void TradeTrackerStarting(this ILogger logger, string symbol)
{
_tradeTrackerStarting(logger, symbol, null);
}
public static void TradeTrackerStartFailed(this ILogger logger, string symbol, string error)
{
_tradeTrackerStartFailed(logger, symbol, error, null);
}
public static void TradeTrackerStarted(this ILogger logger, string symbol)
{
_tradeTrackerStarted(logger, symbol, null);
}
public static void TradeTrackerStopping(this ILogger logger, string symbol)
{
_tradeTrackerStopping(logger, symbol, null);
}
public static void TradeTrackerStopped(this ILogger logger, string symbol)
{
_tradeTrackerStopped(logger, symbol, null);
}
public static void TradeTrackerInitialDataSet(this ILogger logger, string symbol, int count, long lastId)
{
_tradeTrackerInitialDataSet(logger, symbol, count, lastId, null);
}
public static void TradeTrackerPreSnapshotSkip(this ILogger logger, string symbol, long lastId)
{
_tradeTrackerPreSnapshotSkip(logger, symbol, lastId, null);
}
public static void TradeTrackerPreSnapshotApplied(this ILogger logger, string symbol, long lastId)
{
_tradeTrackerPreSnapshotApplied(logger, symbol, lastId, null);
}
public static void TradeTrackerTradeAdded(this ILogger logger, string symbol, long lastId)
{
_tradeTrackerTradeAdded(logger, symbol, lastId, null);
}
public static void TradeTrackerConnectionLost(this ILogger logger, string symbol)
{
_tradeTrackerConnectionLost(logger, symbol, null);
}
public static void TradeTrackerConnectionClosed(this ILogger logger, string symbol)
{
_tradeTrackerConnectionClosed(logger, symbol, null);
}
public static void TradeTrackerConnectionRestored(this ILogger logger, string symbol)
{
_tradeTrackerConnectionRestored(logger, symbol, null);
}
}
}
+27
View File
@@ -68,6 +68,33 @@
Json Json
} }
/// <summary>
/// Tracker sync status
/// </summary>
public enum SyncStatus
{
/// <summary>
/// Not connected
/// </summary>
Disconnected,
/// <summary>
/// Syncing, data connection is being made
/// </summary>
Syncing,
/// <summary>
/// The connection is active, but the full data backlog is not yet reached. For example, a tracker set to retain 10 minutes of data only has 8 minutes of data at this moment.
/// </summary>
PartiallySynced,
/// <summary>
/// Synced
/// </summary>
Synced,
/// <summary>
/// Disposed
/// </summary>
Diposed
}
/// <summary> /// <summary>
/// Status of the order book /// Status of the order book
/// </summary> /// </summary>
@@ -58,12 +58,16 @@ namespace CryptoExchange.Net.Objects
/// </summary> /// </summary>
public IRateLimitGuard? LimitGuard { get; set; } public IRateLimitGuard? LimitGuard { get; set; }
/// <summary> /// <summary>
/// Whether this request should never be cached /// Whether this request should never be cached
/// </summary> /// </summary>
public bool PreventCaching { get; set; } public bool PreventCaching { get; set; }
/// <summary>
/// Connection id
/// </summary>
public int? ConnectionId { get; set; }
/// <summary> /// <summary>
/// ctor /// ctor
/// </summary> /// </summary>
@@ -1,5 +1,6 @@
using CryptoExchange.Net.Interfaces; using CryptoExchange.Net.Interfaces;
using CryptoExchange.Net.Objects.Options; using CryptoExchange.Net.Objects.Options;
using CryptoExchange.Net.SharedApis;
using System; using System;
namespace CryptoExchange.Net.OrderBook namespace CryptoExchange.Net.OrderBook
@@ -8,14 +9,14 @@ namespace CryptoExchange.Net.OrderBook
public class OrderBookFactory<TOptions> : IOrderBookFactory<TOptions> where TOptions: OrderBookOptions public class OrderBookFactory<TOptions> : IOrderBookFactory<TOptions> where TOptions: OrderBookOptions
{ {
private readonly Func<string, Action<TOptions>?, ISymbolOrderBook> _symbolCtor; private readonly Func<string, Action<TOptions>?, ISymbolOrderBook> _symbolCtor;
private readonly Func<string, string, Action<TOptions>?, ISymbolOrderBook> _assetsCtor; private readonly Func<SharedSymbol, Action<TOptions>?, ISymbolOrderBook> _assetsCtor;
/// <summary> /// <summary>
/// ctor /// ctor
/// </summary> /// </summary>
/// <param name="symbolCtor"></param> /// <param name="symbolCtor"></param>
/// <param name="assetsCtor"></param> /// <param name="assetsCtor"></param>
public OrderBookFactory(Func<string, Action<TOptions>?, ISymbolOrderBook> symbolCtor, Func<string, string, Action<TOptions>?, ISymbolOrderBook> assetsCtor) public OrderBookFactory(Func<string, Action<TOptions>?, ISymbolOrderBook> symbolCtor, Func<SharedSymbol, Action<TOptions>?, ISymbolOrderBook> assetsCtor)
{ {
_symbolCtor = symbolCtor; _symbolCtor = symbolCtor;
_assetsCtor = assetsCtor; _assetsCtor = assetsCtor;
@@ -25,6 +26,9 @@ namespace CryptoExchange.Net.OrderBook
public ISymbolOrderBook Create(string symbol, Action<TOptions>? options = null) => _symbolCtor(symbol, options); public ISymbolOrderBook Create(string symbol, Action<TOptions>? options = null) => _symbolCtor(symbol, options);
/// <inheritdoc /> /// <inheritdoc />
public ISymbolOrderBook Create(string baseAsset, string quoteAsset, Action<TOptions>? options = null) => _assetsCtor(baseAsset, quoteAsset, options); public ISymbolOrderBook Create(string baseAsset, string quoteAsset, Action<TOptions>? options = null) => _assetsCtor(new SharedSymbol(TradingMode.Spot, baseAsset, quoteAsset), options);
/// <inheritdoc />
public ISymbolOrderBook Create(SharedSymbol symbol, Action<TOptions>? options = null) => _assetsCtor(symbol, options);
} }
} }
@@ -18,6 +18,10 @@ namespace CryptoExchange.Net.RateLimiting.Guards
/// </summary> /// </summary>
public static Func<RequestDefinition, string, string?, string> PerEndpoint { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => def.Path + def.Method); public static Func<RequestDefinition, string, string?, string> PerEndpoint { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => def.Path + def.Method);
/// <summary> /// <summary>
/// Apply guard per connection
/// </summary>
public static Func<RequestDefinition, string, string?, string> PerConnection { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => def.ConnectionId.ToString());
/// <summary>
/// Apply guard per API key /// Apply guard per API key
/// </summary> /// </summary>
public static Func<RequestDefinition, string, string?, string> PerApiKey { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => key!); public static Func<RequestDefinition, string, string?, string> PerApiKey { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => key!);
@@ -19,7 +19,12 @@ namespace CryptoExchange.Net.Requests
if (client == null) if (client == null)
{ {
var handler = new HttpClientHandler(); var handler = new HttpClientHandler();
handler.AutomaticDecompression = DecompressionMethods.GZip | DecompressionMethods.Deflate; try
{
handler.AutomaticDecompression = DecompressionMethods.GZip | DecompressionMethods.Deflate;
}
catch (PlatformNotSupportedException) { }
if (proxy != null) if (proxy != null)
{ {
handler.Proxy = new WebProxy handler.Proxy = new WebProxy
@@ -18,6 +18,11 @@ namespace CryptoExchange.Net.SharedApis
/// </summary> /// </summary>
TradingMode[] SupportedTradingModes { get; } TradingMode[] SupportedTradingModes { get; }
/// <summary>
/// Whether or not API credentials have been configured for this client. Does not check the credentials are actually valid.
/// </summary>
bool Authenticated { get; }
/// <summary> /// <summary>
/// Format a base and quote asset to an exchange accepted symbol /// Format a base and quote asset to an exchange accepted symbol
/// </summary> /// </summary>
@@ -14,15 +14,15 @@ namespace CryptoExchange.Net.SharedApis
/// <summary> /// <summary>
/// Last trade price /// Last trade price
/// </summary> /// </summary>
public decimal LastPrice { get; set; } public decimal? LastPrice { get; set; }
/// <summary> /// <summary>
/// High price in the last 24h /// High price in the last 24h
/// </summary> /// </summary>
public decimal HighPrice { get; set; } public decimal? HighPrice { get; set; }
/// <summary> /// <summary>
/// Low price in the last 24h /// Low price in the last 24h
/// </summary> /// </summary>
public decimal LowPrice { get; set; } public decimal? LowPrice { get; set; }
/// <summary> /// <summary>
/// The volume in the last 24h /// The volume in the last 24h
/// </summary> /// </summary>
@@ -51,7 +51,7 @@ namespace CryptoExchange.Net.SharedApis
/// <summary> /// <summary>
/// ctor /// ctor
/// </summary> /// </summary>
public SharedFuturesTicker(string symbol, decimal lastPrice, decimal highPrice, decimal lowPrice, decimal volume, decimal? changePercentage) public SharedFuturesTicker(string symbol, decimal? lastPrice, decimal? highPrice, decimal? lowPrice, decimal volume, decimal? changePercentage)
{ {
Symbol = symbol; Symbol = symbol;
LastPrice = lastPrice; LastPrice = lastPrice;
@@ -19,6 +19,10 @@ namespace CryptoExchange.Net.SharedApis
/// Trade time /// Trade time
/// </summary> /// </summary>
public DateTime Timestamp { get; set; } public DateTime Timestamp { get; set; }
/// <summary>
/// Trade side. Buy means that the taker took an ask order of the order book, sell means the taker took a bid order of the order book.
/// </summary>
public SharedOrderSide? Side { get; set; }
/// <summary> /// <summary>
/// ctor /// ctor
@@ -209,7 +209,7 @@ namespace CryptoExchange.Net.Sockets
{ {
if (Parameters.RateLimiter != null) if (Parameters.RateLimiter != null)
{ {
var definition = new RequestDefinition(Id.ToString(), HttpMethod.Get); var definition = new RequestDefinition(Uri.AbsolutePath, HttpMethod.Get) { ConnectionId = Id };
var limitResult = await Parameters.RateLimiter.ProcessAsync(_logger, Id, RateLimitItemType.Connection, definition, _baseAddress, null, 1, Parameters.RateLimitingBehaviour, _ctsSource.Token).ConfigureAwait(false); var limitResult = await Parameters.RateLimiter.ProcessAsync(_logger, Id, RateLimitItemType.Connection, definition, _baseAddress, null, 1, Parameters.RateLimitingBehaviour, _ctsSource.Token).ConfigureAwait(false);
if (!limitResult) if (!limitResult)
return new CallResult(new ClientRateLimitError("Connection limit reached")); return new CallResult(new ClientRateLimitError("Connection limit reached"));
@@ -475,7 +475,7 @@ namespace CryptoExchange.Net.Sockets
/// <returns></returns> /// <returns></returns>
private async Task SendLoopAsync() private async Task SendLoopAsync()
{ {
var requestDefinition = new RequestDefinition(Id.ToString(), HttpMethod.Get); var requestDefinition = new RequestDefinition(Uri.AbsolutePath, HttpMethod.Get) { ConnectionId = Id };
try try
{ {
while (true) while (true)
@@ -14,4 +14,19 @@
/// </summary> /// </summary>
public bool Authenticated { get; set; } public bool Authenticated { get; set; }
} }
/// <summary>
/// Dedicated connection state
/// </summary>
public class DedicatedConnectionState
{
/// <summary>
/// Whether the connection is a dedicated request connection
/// </summary>
public bool IsDedicatedRequestConnection { get; set; }
/// <summary>
/// Whether the dedication request connection should be authenticated
/// </summary>
public bool Authenticated { get; set; }
}
} }
+12 -1
View File
@@ -177,6 +177,10 @@ namespace CryptoExchange.Net.Sockets
/// <inheritdoc /> /// <inheritdoc />
public override async Task<CallResult> Handle(SocketConnection connection, DataEvent<object> message) public override async Task<CallResult> Handle(SocketConnection connection, DataEvent<object> message)
{ {
var typedMessage = message.As((TServerResponse)message.Data);
if (!ValidateMessage(typedMessage))
return new CallResult(null);
CurrentResponses++; CurrentResponses++;
if (CurrentResponses == RequiredResponses) if (CurrentResponses == RequiredResponses)
{ {
@@ -186,7 +190,7 @@ namespace CryptoExchange.Net.Sockets
if (Result?.Success != false) if (Result?.Success != false)
// If an error result is already set don't override that // If an error result is already set don't override that
Result = HandleMessage(connection, message.As((TServerResponse)message.Data)); Result = HandleMessage(connection, typedMessage);
if (CurrentResponses == RequiredResponses) if (CurrentResponses == RequiredResponses)
{ {
@@ -198,6 +202,13 @@ namespace CryptoExchange.Net.Sockets
return Result; return Result;
} }
/// <summary>
/// Validate if a message is actually processable by this query
/// </summary>
/// <param name="message"></param>
/// <returns></returns>
public virtual bool ValidateMessage(DataEvent<TServerResponse> message) => true;
/// <summary> /// <summary>
/// Handle the query response /// Handle the query response
/// </summary> /// </summary>
@@ -186,9 +186,9 @@ namespace CryptoExchange.Net.Sockets
} }
/// <summary> /// <summary>
/// Whether this connection should be kept alive even when there is no subscription /// Info on whether this connection is a dedicated request connection
/// </summary> /// </summary>
public bool DedicatedRequestConnection { get; internal set; } public DedicatedConnectionState DedicatedRequestConnection { get; internal set; } = new DedicatedConnectionState();
private bool _pausedActivity; private bool _pausedActivity;
private readonly object _listenersLock; private readonly object _listenersLock;
@@ -268,7 +268,7 @@ namespace CryptoExchange.Net.Sockets
lock (_listenersLock) lock (_listenersLock)
{ {
foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription)) foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription))
subscription.Confirmed = false; subscription.Reset();
foreach (var query in _listeners.OfType<Query>().ToList()) foreach (var query in _listeners.OfType<Query>().ToList())
{ {
@@ -293,7 +293,7 @@ namespace CryptoExchange.Net.Sockets
lock (_listenersLock) lock (_listenersLock)
{ {
foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription)) foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription))
subscription.Confirmed = false; subscription.Reset();
foreach (var query in _listeners.OfType<Query>().ToList()) foreach (var query in _listeners.OfType<Query>().ToList())
{ {
@@ -618,7 +618,7 @@ namespace CryptoExchange.Net.Sockets
bool shouldCloseConnection; bool shouldCloseConnection;
lock (_listenersLock) lock (_listenersLock)
shouldCloseConnection = _listeners.OfType<Subscription>().All(r => !r.UserSubscription || r.Closed) && !DedicatedRequestConnection; shouldCloseConnection = _listeners.OfType<Subscription>().All(r => !r.UserSubscription || r.Closed) && !DedicatedRequestConnection.IsDedicatedRequestConnection;
if (!anyDuplicateSubscription) if (!anyDuplicateSubscription)
{ {
@@ -841,7 +841,7 @@ namespace CryptoExchange.Net.Sockets
if (!_socket.IsOpen) if (!_socket.IsOpen)
return new CallResult(new WebError("Socket not connected")); return new CallResult(new WebError("Socket not connected"));
if (!DedicatedRequestConnection) if (!DedicatedRequestConnection.IsDedicatedRequestConnection)
{ {
bool anySubscriptions; bool anySubscriptions;
lock (_listenersLock) lock (_listenersLock)
@@ -859,7 +859,7 @@ namespace CryptoExchange.Net.Sockets
lock (_listenersLock) lock (_listenersLock)
{ {
anyAuthenticated = _listeners.OfType<Subscription>().Any(s => s.Authenticated) anyAuthenticated = _listeners.OfType<Subscription>().Any(s => s.Authenticated)
|| (DedicatedRequestConnection && ApiClient.AuthenticationProvider != null); || (DedicatedRequestConnection.IsDedicatedRequestConnection && DedicatedRequestConnection.Authenticated);
} }
if (anyAuthenticated) if (anyAuthenticated)
@@ -130,6 +130,20 @@ namespace CryptoExchange.Net.Sockets
return Task.FromResult(DoHandleMessage(connection, message)); return Task.FromResult(DoHandleMessage(connection, message));
} }
/// <summary>
/// Reset the subscription
/// </summary>
public void Reset()
{
Confirmed = false;
DoHandleReset();
}
/// <summary>
/// Connection has been reset, do any logic for resetting the subscription
/// </summary>
public virtual void DoHandleReset() { }
/// <summary> /// <summary>
/// Handle the update message /// Handle the update message
/// </summary> /// </summary>
@@ -0,0 +1,34 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers
{
/// <summary>
/// Compare value
/// </summary>
public record CompareValue
{
/// <summary>
/// The value difference
/// </summary>
public decimal? Difference { get; set; }
/// <summary>
/// The value difference percentage
/// </summary>
public decimal? PercentageDifference { get; set; }
/// <summary>
/// ctor
/// </summary>
public CompareValue(decimal? value1, decimal? value2)
{
if (value1 == null || value2 == null)
return;
Difference = value2 - value1;
PercentageDifference = value1.Value == 0 ? null : Math.Round(value2.Value / value1.Value * 100 - 100, 4);
}
}
}
@@ -0,0 +1,105 @@
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.SharedApis;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// A tracker for kline data of a symbol
/// </summary>
public interface IKlineTracker
{
/// <summary>
/// The total number of klines
/// </summary>
int Count { get; }
/// <summary>
/// Exchange name
/// </summary>
string Exchange { get; }
/// <summary>
/// Symbol name
/// </summary>
string SymbolName { get; }
/// <summary>
/// Symbol
/// </summary>
SharedSymbol Symbol { get; }
/// <summary>
/// The max number of klines tracked
/// </summary>
int? Limit { get; }
/// <summary>
/// The max age of the data tracked
/// </summary>
TimeSpan? Period { get; }
/// <summary>
/// From which timestamp the trades are registered
/// </summary>
DateTime? SyncedFrom { get; }
/// <summary>
/// Sync status
/// </summary>
SyncStatus Status { get; }
/// <summary>
/// Get the last kline
/// </summary>
SharedKline? Last { get; }
/// <summary>
/// Event for when a new kline is added
/// </summary>
event Func<SharedKline, Task>? OnAdded;
/// <summary>
/// Event for when a kline is removed because it's no longer within the period/limit window
/// </summary>
event Func<SharedKline, Task>? OnRemoved;
/// <summary>
/// Event for when a kline is updated
/// </summary>
event Func<SharedKline, Task> OnUpdated;
/// <summary>
/// Event for when the sync status changes
/// </summary>
event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// Start synchronization
/// </summary>
/// <returns></returns>
Task<CallResult> StartAsync(bool startWithSnapshot = true);
/// <summary>
/// Stop synchronization
/// </summary>
/// <returns></returns>
Task StopAsync();
/// <summary>
/// Get the data tracked
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
IEnumerable<SharedKline> GetData(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
/// <summary>
/// Get statitistics on the klines
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
KlinesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
}
}
@@ -0,0 +1,481 @@
using CryptoExchange.Net.Logging.Extensions;
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.Objects.Sockets;
using CryptoExchange.Net.SharedApis;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <inheritdoc />
public class KlineTracker : IKlineTracker
{
private readonly IKlineSocketClient _socketClient;
private readonly IKlineRestClient _restClient;
private SyncStatus _status;
private bool _startWithSnapshot;
/// <summary>
/// The internal data structure
/// </summary>
protected readonly Dictionary<DateTime, SharedKline> _data = new Dictionary<DateTime, SharedKline>();
/// <summary>
/// The pre-snapshot queue buffering updates received before the snapshot is set and which will be applied after the snapshot was set
/// </summary>
protected readonly List<SharedKline> _preSnapshotQueue = new List<SharedKline>();
/// <summary>
/// Lock for accessing _data
/// </summary>
protected readonly object _lock = new object();
/// <summary>
/// The last time the window was applied
/// </summary>
protected DateTime _lastWindowApplied = DateTime.MinValue;
/// <summary>
/// Whether or not the data has changed since last window was applied
/// </summary>
protected bool _changed = false;
/// <summary>
/// The kline interval
/// </summary>
protected readonly SharedKlineInterval _interval;
/// <summary>
/// Whether the snapshot has been set
/// </summary>
protected bool _snapshotSet;
/// <summary>
/// Logger
/// </summary>
protected readonly ILogger _logger;
/// <summary>
/// Update subscription
/// </summary>
protected UpdateSubscription? _updateSubscription;
/// <summary>
/// The timestamp of the first item
/// </summary>
protected DateTime? _firstTimestamp;
/// <inheritdoc/>
public SyncStatus Status
{
get => _status;
set
{
if (value == _status)
return;
var old = _status;
_status = value;
_logger.KlineTrackerStatusChanged(SymbolName, old, value);
OnStatusChanged?.Invoke(old, _status);
}
}
/// <inheritdoc />
public string Exchange { get; }
/// <inheritdoc />
public string SymbolName { get; }
/// <inheritdoc />
public SharedSymbol Symbol { get; }
/// <inheritdoc/>
public int? Limit { get; }
/// <inheritdoc/>
public TimeSpan? Period { get; }
/// <inheritdoc />
public DateTime? SyncedFrom
{
get
{
if (Period == null)
return _firstTimestamp;
var max = DateTime.UtcNow - Period.Value;
if (_firstTimestamp > max)
return _firstTimestamp;
return max;
}
}
/// <inheritdoc />
public int Count
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.Count;
}
}
}
/// <inheritdoc />
public SharedKline? Last
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.LastOrDefault().Value;
}
}
}
/// <inheritdoc />
public event Func<SharedKline, Task>? OnAdded;
/// <inheritdoc />
public event Func<SharedKline, Task>? OnUpdated;
/// <inheritdoc />
public event Func<SharedKline, Task>? OnRemoved;
/// <inheritdoc />
public event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// ctor
/// </summary>
public KlineTracker(
ILogger? logger,
IKlineRestClient restClient,
IKlineSocketClient socketClient,
SharedSymbol symbol,
SharedKlineInterval interval,
int? limit = null,
TimeSpan? period = null)
{
_logger = logger ?? new NullLogger<KlineTracker>();
Symbol = symbol;
SymbolName = socketClient.FormatSymbol(symbol.BaseAsset, symbol.QuoteAsset, symbol.TradingMode, symbol.DeliverTime);
Exchange = restClient.Exchange;
Limit = limit;
Period = period;
_interval = interval;
_socketClient = socketClient;
_restClient = restClient;
}
/// <inheritdoc />
public async Task<CallResult> StartAsync(bool startWithSnapshot = true)
{
if (Status != SyncStatus.Disconnected)
throw new InvalidOperationException($"Can't start syncing unless state is {SyncStatus.Disconnected}. Current state: {Status}");
_startWithSnapshot = startWithSnapshot;
Status = SyncStatus.Syncing;
_logger.KlineTrackerStarting(SymbolName);
var startResult = await DoStartAsync().ConfigureAwait(false);
if (!startResult)
{
_logger.KlineTrackerStartFailed(SymbolName, startResult.Error!.ToString());
Status = SyncStatus.Disconnected;
return new CallResult(startResult.Error!);
}
_updateSubscription = startResult.Data;
_updateSubscription.ConnectionLost += HandleConnectionLost;
_updateSubscription.ConnectionClosed += HandleConnectionClosed;
_updateSubscription.ConnectionRestored += HandleConnectionRestored;
Status = SyncStatus.Synced;
_logger.KlineTrackerStarted(SymbolName);
return new CallResult(null);
}
/// <inheritdoc />
public async Task StopAsync()
{
_logger.KlineTrackerStopping(SymbolName);
Status = SyncStatus.Disconnected;
await DoStopAsync().ConfigureAwait(false);
_data.Clear();
_preSnapshotQueue.Clear();
_logger.KlineTrackerStopped(SymbolName);
}
/// <summary>
/// The start procedure needed for kline syncing, generally subscribing to an update stream and requesting the snapshot
/// </summary>
/// <returns></returns>
protected virtual async Task<CallResult<UpdateSubscription>> DoStartAsync()
{
var subResult = await _socketClient.SubscribeToKlineUpdatesAsync(new SubscribeKlineRequest(Symbol, _interval),
update =>
{
AddOrUpdate(update.Data);
}).ConfigureAwait(false);
if (!subResult)
{
Status = SyncStatus.Disconnected;
return subResult;
}
if (!_startWithSnapshot)
return subResult;
var startTime = Period == null ? (DateTime?)null : DateTime.UtcNow.Add(-Period.Value);
if (_restClient.GetKlinesOptions.MaxAge != null && DateTime.UtcNow.Add(-_restClient.GetKlinesOptions.MaxAge.Value) > startTime)
startTime = DateTime.UtcNow.Add(-_restClient.GetKlinesOptions.MaxAge.Value);
var limit = Math.Min(_restClient.GetKlinesOptions.MaxRequestDataPoints ?? _restClient.GetKlinesOptions.MaxTotalDataPoints ?? 100, Limit ?? 100);
var request = new GetKlinesRequest(Symbol, _interval, startTime, DateTime.UtcNow, limit: limit);
var data = new List<SharedKline>();
await foreach (var result in ExchangeHelpers.ExecutePages(_restClient.GetKlinesAsync, request).ConfigureAwait(false))
{
if (!result)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(result.Error!);
}
if (Limit != null && data.Count > Limit)
break;
data.AddRange(result.Data);
}
SetInitialData(data);
return subResult;
}
/// <summary>
/// The stop procedure needed, generally stopping the update stream
/// </summary>
/// <returns></returns>
protected virtual Task DoStopAsync() => _updateSubscription?.CloseAsync() ?? Task.CompletedTask;
/// <inheritdoc />
public KlinesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null)
{
var compareTime = SyncedFrom?.AddSeconds(-2);
var stats = GetStats(GetData(fromTimestamp, toTimestamp));
stats.Complete = (fromTimestamp == null || fromTimestamp >= compareTime) && (toTimestamp == null || toTimestamp >= compareTime);
return stats;
}
private KlinesStats GetStats(IEnumerable<SharedKline> klines)
{
if (!klines.Any())
return new KlinesStats();
return new KlinesStats
{
KlineCount = klines.Count(),
FirstOpenTime = klines.First().OpenTime,
LastOpenTime = klines.Last().OpenTime,
HighPrice = klines.Select(d => d.LowPrice).Max(),
LowPrice = klines.Select(d => d.HighPrice).Min(),
Volume = klines.Select(d => d.Volume).Sum(),
AverageVolume = Math.Round(klines.OrderByDescending(d => d.OpenTime).Skip(1).Select(d => d.Volume).DefaultIfEmpty().Average(), 8)
};
}
/// <inheritdoc />
public IEnumerable<SharedKline> GetData(DateTime? since = null, DateTime? until = null)
{
lock (_lock)
{
ApplyWindow(true);
IEnumerable<SharedKline> result = _data.Values;
if (since != null)
result = result.Where(d => d.OpenTime >= since);
if (until != null)
result = result.Where(d => d.OpenTime <= until);
return result.ToList();
}
}
/// <summary>
/// Set the initial kline data snapshot
/// </summary>
/// <param name="data"></param>
protected void SetInitialData(IEnumerable<SharedKline> data)
{
lock (_lock)
{
_data.Clear();
IEnumerable<SharedKline> items = data.OrderByDescending(d => d.OpenTime);
if (Limit != null)
items = items.Take(Limit.Value);
if (Period != null)
items = items.Where(e => e.OpenTime >= DateTime.UtcNow.Add(-Period.Value));
foreach (var item in items.OrderBy(d => d.OpenTime))
_data.Add(item.OpenTime, item);
_snapshotSet = true;
foreach (var item in _preSnapshotQueue)
{
if (_data.ContainsKey(item.OpenTime))
continue;
_data.Add(item.OpenTime, item);
}
_firstTimestamp = _data.Min(v => v.Key);
ApplyWindow(false);
_logger.KlineTrackerInitialDataSet(SymbolName, _data.Last().Key);
}
}
/// <summary>
/// Add or update a kline
/// </summary>
/// <param name="item"></param>
protected void AddOrUpdate(SharedKline item) => AddOrUpdate(new[] { item });
/// <summary>
/// Add or update klines
/// </summary>
/// <param name="items"></param>
protected void AddOrUpdate(IEnumerable<SharedKline> items)
{
lock (_lock)
{
if (_restClient != null && _startWithSnapshot && !_snapshotSet)
{
_preSnapshotQueue.AddRange(items);
return;
}
foreach (var item in items)
{
if (_data.TryGetValue(item.OpenTime, out var existing))
{
_data.Remove(item.OpenTime);
_data.Add(item.OpenTime, item);
OnUpdated?.Invoke(item);
_logger.KlineTrackerKlineUpdated(SymbolName, _data.Last().Key);
}
else
{
_data.Add(item.OpenTime, item);
OnAdded?.Invoke(item);
_logger.KlineTrackerKlineAdded(SymbolName, _data.Last().Key);
}
}
_firstTimestamp = _data.Min(x => x.Key);
_changed = true;
SetSyncStatus();
ApplyWindow(true);
}
}
private void ApplyWindow(bool broadcastEvents)
{
if (!_changed && (DateTime.UtcNow - _lastWindowApplied) < TimeSpan.FromSeconds(1))
return;
if (Period != null)
{
var compareDate = DateTime.UtcNow.Add(-Period.Value);
for (var i = 0; i < _data.Count; i++)
{
var item = _data.ElementAt(0);
if (item.Key >= compareDate)
break;
_data.Remove(item.Key);
if (broadcastEvents)
OnRemoved?.Invoke(item.Value);
}
}
if (Limit != null && _data.Count > Limit.Value)
{
var toRemove = Math.Max(0, _data.Count - Limit.Value);
for (var i = 0; i < toRemove; i++)
{
var item = _data.ElementAt(0);
_data.Remove(item.Key);
if (broadcastEvents)
OnRemoved?.Invoke(item.Value);
}
}
_lastWindowApplied = DateTime.UtcNow;
_changed = false;
}
private void HandleConnectionLost()
{
_logger.KlineTrackerConnectionLost(SymbolName);
if (Status != SyncStatus.Disconnected)
{
Status = SyncStatus.Syncing;
_snapshotSet = false;
_firstTimestamp = null;
_preSnapshotQueue.Clear();
}
}
private void HandleConnectionClosed()
{
_logger.KlineTrackerConnectionClosed(SymbolName);
Status = SyncStatus.Disconnected;
_ = StopAsync();
}
private async void HandleConnectionRestored(TimeSpan _)
{
Status = SyncStatus.Syncing;
var success = false;
while (!success)
{
if (Status != SyncStatus.Syncing)
return;
var resyncResult = await DoStartAsync().ConfigureAwait(false);
success = resyncResult;
}
_logger.KlineTrackerConnectionRestored(SymbolName);
SetSyncStatus();
}
private void SetSyncStatus()
{
if (Status == SyncStatus.Synced)
return;
if (Period != null)
{
if (_firstTimestamp <= DateTime.UtcNow - Period.Value)
Status = SyncStatus.Synced;
else
Status = SyncStatus.PartiallySynced;
}
if (Limit != null)
{
if (_data.Count == Limit.Value)
Status = SyncStatus.Synced;
else
Status = SyncStatus.PartiallySynced;
}
if (Period == null && Limit == null)
Status = SyncStatus.Synced;
}
}
}
@@ -0,0 +1,30 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// Klines statistics comparison
/// </summary>
public record KlinesCompare
{
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? LowPriceDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? HighPriceDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? VolumeDif { get; set; }
/// <summary>
/// Number of trades
/// </summary>
public CompareValue? AverageVolumeDif { get; set; }
}
}
@@ -0,0 +1,59 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Klines
{
/// <summary>
/// Klines statistics
/// </summary>
public record KlinesStats
{
/// <summary>
/// Number of klines
/// </summary>
public int KlineCount { get; set; }
/// <summary>
/// The kline open time of the first entry
/// </summary>
public DateTime? FirstOpenTime { get; set; }
/// <summary>
/// The kline open time of the last entry
/// </summary>
public DateTime? LastOpenTime { get; set; }
/// <summary>
/// Lowest trade price
/// </summary>
public decimal? LowPrice { get; set; }
/// <summary>
/// Highest trade price
/// </summary>
public decimal? HighPrice { get; set; }
/// <summary>
/// Trade volume
/// </summary>
public decimal Volume { get; set; }
/// <summary>
/// Average volume per kline
/// </summary>
public decimal? AverageVolume { get; set; }
/// <summary>
/// Whether the data is complete
/// </summary>
public bool Complete { get; set; }
/// <summary>
/// Compare 2 stat snapshots to eachother
/// </summary>
public KlinesCompare CompareTo(KlinesStats otherStats)
{
return new KlinesCompare
{
LowPriceDif = new CompareValue(LowPrice, otherStats.LowPrice),
HighPriceDif = new CompareValue(HighPrice, otherStats.HighPrice),
VolumeDif = new CompareValue(Volume, otherStats.Volume),
AverageVolumeDif = new CompareValue(AverageVolume, otherStats.AverageVolume),
};
}
}
}
@@ -0,0 +1,100 @@
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.SharedApis;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// A tracker for trades on a symbol
/// </summary>
public interface ITradeTracker
{
/// <summary>
/// The total number of trades
/// </summary>
int Count { get; }
/// <summary>
/// Exchange name
/// </summary>
string Exchange { get; }
/// <summary>
/// Symbol name
/// </summary>
string SymbolName { get; }
/// <summary>
/// Symbol
/// </summary>
SharedSymbol Symbol { get; }
/// <summary>
/// The max number of trades tracked
/// </summary>
int? Limit { get; }
/// <summary>
/// The max age of the data tracked
/// </summary>
TimeSpan? Period { get; }
/// <summary>
/// From which timestamp the trades are registered
/// </summary>
DateTime? SyncedFrom { get; }
/// <summary>
/// The current synchronization status
/// </summary>
SyncStatus Status { get; }
/// <summary>
/// Get the last trade
/// </summary>
SharedTrade? Last { get; }
/// <summary>
/// Event for when a new trade is added
/// </summary>
event Func<SharedTrade, Task>? OnAdded;
/// <summary>
/// Event for when a trade is removed because it's no longer within the period/limit window
/// </summary>
event Func<SharedTrade, Task>? OnRemoved;
/// <summary>
/// Event for when the sync status changes
/// </summary>
event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// Start synchronization
/// </summary>
/// <returns></returns>
Task<CallResult> StartAsync(bool startWithSnapshot = true);
/// <summary>
/// Stop synchronization
/// </summary>
/// <returns></returns>
Task StopAsync();
/// <summary>
/// Get the data tracked
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
IEnumerable<SharedTrade> GetData(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
/// <summary>
/// Get statitistics on the trades
/// </summary>
/// <param name="fromTimestamp">Start timestamp to get the data from, defaults to tracked data start time</param>
/// <param name="toTimestamp">End timestamp to get the data until, defaults to current time</param>
/// <returns></returns>
TradesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null);
}
}
@@ -0,0 +1,495 @@
using CryptoExchange.Net.Logging.Extensions;
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.Objects.Sockets;
using CryptoExchange.Net.SharedApis;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <inheritdoc />
public class TradeTracker : ITradeTracker
{
private readonly ITradeSocketClient _socketClient;
private readonly IRecentTradeRestClient? _recentRestClient;
private readonly ITradeHistoryRestClient? _historyRestClient;
private SyncStatus _status;
private long _snapshotId;
private bool _startWithSnapshot;
/// <summary>
/// The internal data structure
/// </summary>
protected readonly List<SharedTrade> _data = new List<SharedTrade>();
/// <summary>
/// The pre-snapshot queue buffering updates received before the snapshot is set and which will be applied after the snapshot was set
/// </summary>
protected readonly List<SharedTrade> _preSnapshotQueue = new List<SharedTrade>();
/// <summary>
/// The last time the window was applied
/// </summary>
protected DateTime _lastWindowApplied = DateTime.MinValue;
/// <summary>
/// Whether or not the data has changed since last window was applied
/// </summary>
protected bool _changed = false;
/// <summary>
/// Lock for accessing _data
/// </summary>
protected readonly object _lock = new object();
/// <summary>
/// Whether the snapshot has been set
/// </summary>
protected bool _snapshotSet;
/// <summary>
/// Logger
/// </summary>
protected readonly ILogger _logger;
/// <summary>
/// Update subscription
/// </summary>
protected UpdateSubscription? _updateSubscription;
/// <summary>
/// The timestamp of the first item
/// </summary>
protected DateTime? _firstTimestamp;
/// <inheritdoc />
public string Exchange { get; }
/// <inheritdoc />
public string SymbolName { get; }
/// <inheritdoc />
public SharedSymbol Symbol { get; }
/// <inheritdoc/>
public int? Limit { get; }
/// <inheritdoc/>
public TimeSpan? Period { get; }
/// <inheritdoc/>
public SyncStatus Status
{
get => _status;
set
{
if (value == _status)
return;
var old = _status;
_status = value;
_logger.TradeTrackerStatusChanged(SymbolName, old, value);
OnStatusChanged?.Invoke(old, _status);
}
}
/// <inheritdoc />
public int Count
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.Count;
}
}
}
/// <inheritdoc />
public DateTime? SyncedFrom
{
get
{
if (Period == null)
return _firstTimestamp;
var max = DateTime.UtcNow - Period.Value;
if (_firstTimestamp > max)
return _firstTimestamp;
return max;
}
}
/// <inheritdoc />
public SharedTrade? Last
{
get
{
lock (_lock)
{
ApplyWindow(true);
return _data.LastOrDefault();
}
}
}
/// <inheritdoc />
public event Func<SharedTrade, Task>? OnAdded;
/// <inheritdoc />
public event Func<SharedTrade, Task>? OnRemoved;
/// <inheritdoc />
public event Func<SyncStatus, SyncStatus, Task>? OnStatusChanged;
/// <summary>
/// ctor
/// </summary>
public TradeTracker(
ILogger? logger,
IRecentTradeRestClient? recentRestClient,
ITradeHistoryRestClient? historyRestClient,
ITradeSocketClient socketClient,
SharedSymbol symbol,
int? limit = null,
TimeSpan? period = null)
{
_logger = logger ?? new NullLogger<TradeTracker>();
_recentRestClient = recentRestClient;
_historyRestClient = historyRestClient;
_socketClient = socketClient;
Exchange = socketClient.Exchange;
Symbol = symbol;
SymbolName = socketClient.FormatSymbol(symbol.BaseAsset, symbol.QuoteAsset, symbol.TradingMode, symbol.DeliverTime);
Limit = limit;
Period = period;
}
private TradesStats GetStats(IEnumerable<SharedTrade> trades)
{
if (!trades.Any())
return new TradesStats();
return new TradesStats
{
TradeCount = trades.Count(),
FirstTradeTime = trades.First().Timestamp,
LastTradeTime = trades.Last().Timestamp,
AveragePrice = Math.Round(trades.Select(d => d.Price).DefaultIfEmpty().Average(), 8),
VolumeWeightedAveragePrice = trades.Any() ? Math.Round(trades.Select(d => d.Price * d.Quantity).DefaultIfEmpty().Sum() / trades.Select(d => d.Quantity).DefaultIfEmpty().Sum(), 8) : null,
Volume = Math.Round(trades.Sum(d => d.Quantity), 8),
QuoteVolume = Math.Round(trades.Sum(d => d.Quantity * d.Price), 8),
BuySellRatio = Math.Round(trades.Where(x => x.Side == SharedOrderSide.Buy).Sum(x => x.Quantity) / trades.Sum(x => x.Quantity), 8)
};
}
/// <inheritdoc />
public TradesStats GetStats(DateTime? fromTimestamp = null, DateTime? toTimestamp = null)
{
var compareTime = SyncedFrom?.AddSeconds(-2);
var stats = GetStats(GetData(fromTimestamp, toTimestamp));
stats.Complete = (fromTimestamp == null || fromTimestamp >= compareTime) && (toTimestamp == null || toTimestamp >= compareTime);
return stats;
}
/// <inheritdoc />
public async Task<CallResult> StartAsync(bool startWithSnapshot = true)
{
if (Status != SyncStatus.Disconnected)
throw new InvalidOperationException($"Can't start syncing unless state is {SyncStatus.Disconnected}. Current state: {Status}");
_startWithSnapshot = startWithSnapshot;
Status = SyncStatus.Syncing;
_logger.TradeTrackerStarting(SymbolName);
var subResult = await DoStartAsync().ConfigureAwait(false);
if (!subResult)
{
_logger.TradeTrackerStartFailed(SymbolName, subResult.Error!.ToString());
Status = SyncStatus.Disconnected;
return subResult;
}
_updateSubscription = subResult.Data;
_updateSubscription.ConnectionLost += HandleConnectionLost;
_updateSubscription.ConnectionClosed += HandleConnectionClosed;
_updateSubscription.ConnectionRestored += HandleConnectionRestored;
SetSyncStatus();
_logger.TradeTrackerStarted(SymbolName);
return new CallResult(null);
}
/// <inheritdoc />
public async Task StopAsync()
{
_logger.TradeTrackerStopping(SymbolName);
Status = SyncStatus.Disconnected;
await DoStopAsync().ConfigureAwait(false);
_data.Clear();
_preSnapshotQueue.Clear();
_logger.TradeTrackerStopped(SymbolName);
}
/// <summary>
/// The start procedure needed for trade syncing, generally subscribing to an update stream and requesting the snapshot
/// </summary>
/// <returns></returns>
protected virtual async Task<CallResult<UpdateSubscription>> DoStartAsync()
{
var subResult = await _socketClient.SubscribeToTradeUpdatesAsync(new SubscribeTradeRequest(Symbol),
update =>
{
AddData(update.Data);
}).ConfigureAwait(false);
if (!subResult)
{
Status = SyncStatus.Disconnected;
return subResult;
}
if (!_startWithSnapshot)
return subResult;
if (_historyRestClient != null)
{
var startTime = Period == null ? DateTime.UtcNow.AddMinutes(-5) : DateTime.UtcNow.Add(-Period.Value);
var request = new GetTradeHistoryRequest(Symbol, startTime, DateTime.UtcNow);
var data = new List<SharedTrade>();
await foreach(var result in ExchangeHelpers.ExecutePages(_historyRestClient.GetTradeHistoryAsync, request).ConfigureAwait(false))
{
if (!result)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(result.Error!);
}
if (Limit != null && data.Count > Limit)
break;
data.AddRange(result.Data);
}
SetInitialData(data);
}
else if (_recentRestClient != null)
{
int? limit = null;
if (Limit.HasValue)
limit = Math.Min(_recentRestClient.GetRecentTradesOptions.MaxLimit, Limit.Value);
var snapshot = await _recentRestClient.GetRecentTradesAsync(new GetRecentTradesRequest(Symbol, limit)).ConfigureAwait(false);
if (!snapshot)
{
_ = subResult.Data.CloseAsync();
Status = SyncStatus.Disconnected;
return subResult.AsError<UpdateSubscription>(snapshot.Error!);
}
SetInitialData(snapshot.Data);
}
return subResult;
}
/// <summary>
/// The stop procedure needed, generally stopping the update stream
/// </summary>
/// <returns></returns>
protected virtual Task DoStopAsync() => _updateSubscription?.CloseAsync() ?? Task.CompletedTask;
/// <inheritdoc />
public IEnumerable<SharedTrade> GetData(DateTime? since = null, DateTime? until = null)
{
lock (_lock)
{
ApplyWindow(true);
IEnumerable<SharedTrade> result = _data;
if (since != null)
result = result.Where(d => d.Timestamp >= since);
if (until != null)
result = result.Where(d => d.Timestamp <= until);
return result.ToList();
}
}
/// <summary>
/// Set the initial trade data snapshot
/// </summary>
/// <param name="data"></param>
protected void SetInitialData(IEnumerable<SharedTrade> data)
{
lock (_lock)
{
_data.Clear();
IEnumerable<SharedTrade> items = data.OrderByDescending(d => d.Timestamp);
if (Limit != null)
items = items.Take(Limit.Value);
if (Period != null)
items = items.Where(e => e.Timestamp >= DateTime.UtcNow.Add(-Period.Value));
_snapshotId = data.Max(d => d.Timestamp.Ticks);
foreach (var item in items.OrderBy(d => d.Timestamp))
_data.Add(item);
_snapshotSet = true;
_changed = true;
_logger.TradeTrackerInitialDataSet(SymbolName, _data.Count, _snapshotId);
foreach (var item in _preSnapshotQueue)
{
if (_snapshotId >= item.Timestamp.Ticks)
{
_logger.TradeTrackerPreSnapshotSkip(SymbolName, item.Timestamp.Ticks);
continue;
}
_logger.TradeTrackerPreSnapshotApplied(SymbolName, item.Timestamp.Ticks);
_data.Add(item);
}
_firstTimestamp = _data.Min(v => v.Timestamp);
ApplyWindow(false);
}
}
/// <summary>
/// Add a trade
/// </summary>
/// <param name="item"></param>
protected void AddData(SharedTrade item) => AddData(new[] { item });
/// <summary>
/// Add a list of trades
/// </summary>
/// <param name="items"></param>
protected void AddData(IEnumerable<SharedTrade> items)
{
lock (_lock)
{
if ((_recentRestClient != null || _historyRestClient != null) && _startWithSnapshot && !_snapshotSet)
{
_preSnapshotQueue.AddRange(items);
return;
}
foreach (var item in items)
{
_logger.TradeTrackerTradeAdded(SymbolName, item.Timestamp.Ticks);
_data.Add(item);
OnAdded?.Invoke(item);
}
_firstTimestamp = _data.Min(x => x.Timestamp);
_changed = true;
SetSyncStatus();
ApplyWindow(true);
}
}
private void ApplyWindow(bool broadcastEvents)
{
if (!_changed && (DateTime.UtcNow - _lastWindowApplied) < TimeSpan.FromSeconds(1))
return;
if (Period != null)
{
var compareDate = DateTime.UtcNow.Add(-Period.Value);
for(var i = 0; i < _data.Count; i++)
{
var item = _data[0];
if (item.Timestamp >= compareDate)
break;
_data.Remove(item);
if (broadcastEvents)
OnRemoved?.Invoke(item);
}
}
if (Limit != null && _data.Count > Limit.Value)
{
var toRemove = _data.Count - Limit.Value;
for (var i = 0; i < toRemove; i++)
{
var item = _data[0];
_data.Remove(item);
if (broadcastEvents)
OnRemoved?.Invoke(item);
}
}
_lastWindowApplied = DateTime.UtcNow;
_changed = false;
if (Status == SyncStatus.PartiallySynced)
// Need to check if sync status should be changed even if there may not be any new data
SetSyncStatus();
}
private void HandleConnectionLost()
{
_logger.TradeTrackerConnectionLost(SymbolName);
if (Status != SyncStatus.Disconnected)
{
Status = SyncStatus.Syncing;
_snapshotSet = false;
_firstTimestamp = null;
_preSnapshotQueue.Clear();
}
}
private void HandleConnectionClosed()
{
_logger.TradeTrackerConnectionClosed(SymbolName);
Status = SyncStatus.Disconnected;
_ = StopAsync();
}
private async void HandleConnectionRestored(TimeSpan _)
{
Status = SyncStatus.Syncing;
var success = false;
while (!success)
{
if (Status != SyncStatus.Syncing)
return;
var resyncResult = await DoStartAsync().ConfigureAwait(false);
success = resyncResult;
}
_logger.TradeTrackerConnectionRestored(SymbolName);
SetSyncStatus();
}
private void SetSyncStatus()
{
if (Status == SyncStatus.Synced)
return;
if (Period != null)
{
if (_firstTimestamp <= DateTime.UtcNow - Period.Value)
Status = SyncStatus.Synced;
else
Status = SyncStatus.PartiallySynced;
}
if (Limit != null)
{
if (_data.Count == Limit.Value)
Status = SyncStatus.Synced;
else
Status = SyncStatus.PartiallySynced;
}
if (Period == null && Limit == null)
Status = SyncStatus.Synced;
}
}
}
@@ -0,0 +1,37 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// Trades statistics comparison
/// </summary>
public record TradesCompare
{
/// <summary>
/// Number of trades
/// </summary>
public CompareValue TradeCountDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// Average trade price
/// </summary>
public CompareValue? AveragePriceDif { get; set; }
/// <summary>
/// Volume weighted average trade price
/// </summary>
public CompareValue? VolumeWeightedAveragePriceDif { get; set; }
/// <summary>
/// Volume of the trades
/// </summary>
public CompareValue VolumeDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// Volume of the trades in quote asset
/// </summary>
public CompareValue QuoteVolumeDif { get; set; } = new CompareValue(null, null);
/// <summary>
/// The volume weighted Buy/Sell ratio. A 0.7 ratio means 70% of the trade volume was a buy.
/// </summary>
public CompareValue? BuySellRatioDif { get; set; }
}
}
@@ -0,0 +1,65 @@
using System;
using System.Collections.Generic;
using System.Text;
namespace CryptoExchange.Net.Trackers.Trades
{
/// <summary>
/// Trades statistics
/// </summary>
public record TradesStats
{
/// <summary>
/// Number of trades
/// </summary>
public int TradeCount { get; set; }
/// <summary>
/// Timestamp of the last trade
/// </summary>
public DateTime? FirstTradeTime { get; set; }
/// <summary>
/// Timestamp of the first trade
/// </summary>
public DateTime? LastTradeTime { get; set; }
/// <summary>
/// Average trade price
/// </summary>
public decimal? AveragePrice { get; set; }
/// <summary>
/// Volume weighted average trade price
/// </summary>
public decimal? VolumeWeightedAveragePrice { get; set; }
/// <summary>
/// Volume of the trades
/// </summary>
public decimal Volume { get; set; }
/// <summary>
/// Volume of the trades in quote asset
/// </summary>
public decimal QuoteVolume { get; set; }
/// <summary>
/// The volume weighted Buy/Sell ratio. A 0.7 ratio means 70% of the trade volume was a buy.
/// </summary>
public decimal? BuySellRatio { get; set; }
/// <summary>
/// Whether the data is complete
/// </summary>
public bool Complete { get; set; }
/// <summary>
/// Compare 2 stat snapshots to eachother
/// </summary>
public TradesCompare CompareTo(TradesStats otherStats)
{
return new TradesCompare
{
TradeCountDif = new CompareValue(TradeCount, otherStats.TradeCount),
AveragePriceDif = new CompareValue(AveragePrice, otherStats.AveragePrice),
VolumeWeightedAveragePriceDif = new CompareValue(VolumeWeightedAveragePrice, otherStats.VolumeWeightedAveragePrice),
VolumeDif = new CompareValue(Volume, otherStats.Volume),
QuoteVolumeDif = new CompareValue(QuoteVolume, otherStats.QuoteVolume),
BuySellRatioDif = new CompareValue(BuySellRatio, otherStats.BuySellRatio),
};
}
}
}
+16 -15
View File
@@ -5,21 +5,22 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="10.5.0" /> <PackageReference Include="Binance.Net" Version="10.8.0" />
<PackageReference Include="Bitfinex.Net" Version="7.8.0" /> <PackageReference Include="Bitfinex.Net" Version="7.9.0" />
<PackageReference Include="BitMart.Net" Version="1.2.0" /> <PackageReference Include="BitMart.Net" Version="1.5.0" />
<PackageReference Include="Bybit.Net" Version="3.14.0" /> <PackageReference Include="Bybit.Net" Version="3.15.0" />
<PackageReference Include="CoinEx.Net" Version="7.7.0" /> <PackageReference Include="CoinEx.Net" Version="7.8.0" />
<PackageReference Include="GateIo.Net" Version="1.6.0" /> <PackageReference Include="CryptoCom.Net" Version="1.1.0" />
<PackageReference Include="JK.BingX.Net" Version="1.11.0" /> <PackageReference Include="GateIo.Net" Version="1.10.0" />
<PackageReference Include="JK.Bitget.Net" Version="1.10.0" /> <PackageReference Include="JK.BingX.Net" Version="1.12.0" />
<PackageReference Include="JK.Mexc.Net" Version="1.8.0" /> <PackageReference Include="JK.Bitget.Net" Version="1.11.0" />
<PackageReference Include="JK.OKX.Net" Version="2.4.0" /> <PackageReference Include="JK.Mexc.Net" Version="1.10.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="1.0.0" /> <PackageReference Include="JK.OKX.Net" Version="2.7.0" />
<PackageReference Include="JKorf.HTX.Net" Version="6.1.0" /> <PackageReference Include="JKorf.Coinbase.Net" Version="1.2.0" />
<PackageReference Include="KrakenExchange.Net" Version="4.12.0" /> <PackageReference Include="JKorf.HTX.Net" Version="6.3.0" />
<PackageReference Include="Kucoin.Net" Version="5.14.0" /> <PackageReference Include="KrakenExchange.Net" Version="5.1.0" />
<PackageReference Include="Serilog.AspNetCore" Version="8.0.2" /> <PackageReference Include="Kucoin.Net" Version="5.17.0" />
<PackageReference Include="Serilog.AspNetCore" Version="8.0.3" />
</ItemGroup> </ItemGroup>
</Project> </Project>
+5
View File
@@ -7,6 +7,7 @@
@inject IBybitRestClient bybitClient @inject IBybitRestClient bybitClient
@inject ICoinbaseRestClient coinbaseClient @inject ICoinbaseRestClient coinbaseClient
@inject ICoinExRestClient coinexClient @inject ICoinExRestClient coinexClient
@inject ICryptoComRestClient cryptocomClient
@inject IGateIoRestClient gateioClient @inject IGateIoRestClient gateioClient
@inject IHTXRestClient huobiClient @inject IHTXRestClient huobiClient
@inject IKrakenRestClient krakenClient @inject IKrakenRestClient krakenClient
@@ -33,6 +34,7 @@
var bybitTask = bybitClient.V5Api.ExchangeData.GetSpotTickersAsync("BTCUSDT"); var bybitTask = bybitClient.V5Api.ExchangeData.GetSpotTickersAsync("BTCUSDT");
var coinbaseTask = coinbaseClient.AdvancedTradeApi.ExchangeData.GetSymbolAsync("BTC-USDT"); var coinbaseTask = coinbaseClient.AdvancedTradeApi.ExchangeData.GetSymbolAsync("BTC-USDT");
var coinexTask = coinexClient.SpotApi.ExchangeData.GetTickerAsync("BTCUSDT"); var coinexTask = coinexClient.SpotApi.ExchangeData.GetTickerAsync("BTCUSDT");
var cryptocomTask = cryptocomClient.ExchangeApi.ExchangeData.GetTickersAsync("BTC_USDT");
var gateioTask = gateioClient.SpotApi.ExchangeData.GetTickersAsync("BTC_USDT"); var gateioTask = gateioClient.SpotApi.ExchangeData.GetTickersAsync("BTC_USDT");
var htxTask = huobiClient.SpotApi.ExchangeData.GetTickerAsync("btcusdt"); var htxTask = huobiClient.SpotApi.ExchangeData.GetTickerAsync("btcusdt");
var krakenTask = krakenClient.SpotApi.ExchangeData.GetTickerAsync("XBTUSD"); var krakenTask = krakenClient.SpotApi.ExchangeData.GetTickerAsync("XBTUSD");
@@ -66,6 +68,9 @@
if (coinexTask.Result.Success) if (coinexTask.Result.Success)
_prices.Add("CoinEx", coinexTask.Result.Data.Ticker.LastPrice); _prices.Add("CoinEx", coinexTask.Result.Data.Ticker.LastPrice);
if (cryptocomTask.Result.Success)
_prices.Add("CryptoCom", cryptocomTask.Result.Data.First().LastPrice ?? 0);
if (gateioTask.Result.Success) if (gateioTask.Result.Success)
_prices.Add("GateIo", gateioTask.Result.Data.First().LastPrice); _prices.Add("GateIo", gateioTask.Result.Data.First().LastPrice);
+3 -1
View File
@@ -7,6 +7,7 @@
@inject IBybitSocketClient bybitSocketClient @inject IBybitSocketClient bybitSocketClient
@inject ICoinbaseSocketClient coinbaseSocketClient @inject ICoinbaseSocketClient coinbaseSocketClient
@inject ICoinExSocketClient coinExSocketClient @inject ICoinExSocketClient coinExSocketClient
@inject ICryptoComSocketClient cryptocomSocketClient
@inject IGateIoSocketClient gateioSocketClient @inject IGateIoSocketClient gateioSocketClient
@inject IHTXSocketClient htxSocketClient @inject IHTXSocketClient htxSocketClient
@inject IKrakenSocketClient krakenSocketClient @inject IKrakenSocketClient krakenSocketClient
@@ -41,9 +42,10 @@
bybitSocketClient.V5SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bybit", data.Data.LastPrice)), bybitSocketClient.V5SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bybit", data.Data.LastPrice)),
coinExSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("CoinEx", data.Data.LastPrice)), coinExSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("CoinEx", data.Data.LastPrice)),
coinbaseSocketClient.AdvancedTradeApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Coinbase", data.Data.LastPrice)), coinbaseSocketClient.AdvancedTradeApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Coinbase", data.Data.LastPrice)),
cryptocomSocketClient.ExchangeApi.SubscribeToTickerUpdatesAsync("ETH_BTC", data => UpdateData("CryptoCom", data.Data.LastPrice ?? 0)),
gateioSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH_BTC", data => UpdateData("GateIo", data.Data.LastPrice)), gateioSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH_BTC", data => UpdateData("GateIo", data.Data.LastPrice)),
htxSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ethbtc", data => UpdateData("HTX", data.Data.ClosePrice ?? 0)), htxSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ethbtc", data => UpdateData("HTX", data.Data.ClosePrice ?? 0)),
krakenSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH/XBT", data => UpdateData("Kraken", data.Data.LastTrade.Price)), krakenSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH/XBT", data => UpdateData("Kraken", data.Data.LastPrice)),
kucoinSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Kucoin", data.Data.LastPrice ?? 0)), kucoinSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Kucoin", data.Data.LastPrice ?? 0)),
mexcSocketClient.SpotApi.SubscribeToMiniTickerUpdatesAsync("ETHBTC", data => UpdateData("Mexc", data.Data.LastPrice)), mexcSocketClient.SpotApi.SubscribeToMiniTickerUpdatesAsync("ETHBTC", data => UpdateData("Mexc", data.Data.LastPrice)),
okxSocketClient.UnifiedApi.ExchangeData.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("OKX", data.Data.LastPrice ?? 0)), okxSocketClient.UnifiedApi.ExchangeData.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("OKX", data.Data.LastPrice ?? 0)),
@@ -10,6 +10,7 @@
@using CoinEx.Net.Interfaces @using CoinEx.Net.Interfaces
@using Coinbase.Net.Interfaces @using Coinbase.Net.Interfaces
@using CryptoExchange.Net.Interfaces @using CryptoExchange.Net.Interfaces
@using CryptoCom.Net.Interfaces
@using GateIo.Net.Interfaces @using GateIo.Net.Interfaces
@using HTX.Net.Interfaces @using HTX.Net.Interfaces
@using Kraken.Net.Interfaces @using Kraken.Net.Interfaces
@@ -25,6 +26,7 @@
@inject IBybitOrderBookFactory bybitFactory @inject IBybitOrderBookFactory bybitFactory
@inject ICoinbaseOrderBookFactory coinbaseFactory @inject ICoinbaseOrderBookFactory coinbaseFactory
@inject ICoinExOrderBookFactory coinExFactory @inject ICoinExOrderBookFactory coinExFactory
@inject ICryptoComOrderBookFactory cryptocomFactory
@inject IGateIoOrderBookFactory gateioFactory @inject IGateIoOrderBookFactory gateioFactory
@inject IHTXOrderBookFactory htxFactory @inject IHTXOrderBookFactory htxFactory
@inject IKrakenOrderBookFactory krakenFactory @inject IKrakenOrderBookFactory krakenFactory
@@ -72,6 +74,7 @@
{ "Bybit", bybitFactory.Create("ETHBTC", Bybit.Net.Enums.Category.Spot) }, { "Bybit", bybitFactory.Create("ETHBTC", Bybit.Net.Enums.Category.Spot) },
{ "Coinbase", coinbaseFactory.Create("ETH-BTC", null) }, { "Coinbase", coinbaseFactory.Create("ETH-BTC", null) },
{ "CoinEx", coinExFactory.CreateSpot("ETHBTC") }, { "CoinEx", coinExFactory.CreateSpot("ETHBTC") },
{ "CryptoCom", cryptocomFactory.Create("ETH_BTC") },
{ "GateIo", gateioFactory.CreateSpot("ETH_BTC") }, { "GateIo", gateioFactory.CreateSpot("ETH_BTC") },
{ "HTX", htxFactory.CreateSpot("ethbtc") }, { "HTX", htxFactory.CreateSpot("ethbtc") },
{ "Kraken", krakenFactory.CreateSpot("ETH/XBT") }, { "Kraken", krakenFactory.CreateSpot("ETH/XBT") },
+111
View File
@@ -0,0 +1,111 @@
@page "/Trackers"
@using System.Collections.Concurrent
@using System.Timers
@using Binance.Net.Interfaces
@using BingX.Net.Interfaces
@using Bitfinex.Net.Interfaces
@using Bitget.Net.Interfaces;
@using BitMart.Net.Interfaces;
@using Bybit.Net.Interfaces
@using CoinEx.Net.Interfaces
@using Coinbase.Net.Interfaces
@using CryptoExchange.Net.Interfaces
@using CryptoCom.Net.Interfaces
@using CryptoExchange.Net.SharedApis
@using CryptoExchange.Net.Trackers.Trades
@using GateIo.Net.Interfaces
@using HTX.Net.Interfaces
@using Kraken.Net.Interfaces
@using Kucoin.Net.Clients
@using Kucoin.Net.Interfaces
@using Mexc.Net.Interfaces
@using OKX.Net.Interfaces;
@inject IBinanceTrackerFactory binanceFactory
@inject IBingXTrackerFactory bingXFactory
@inject IBitfinexTrackerFactory bitfinexFactory
@inject IBitgetTrackerFactory bitgetFactory
@inject IBitMartTrackerFactory bitmartFactory
@inject IBybitTrackerFactory bybitFactory
@inject ICoinbaseTrackerFactory coinbaseFactory
@inject ICoinExTrackerFactory coinExFactory
@inject ICryptoComTrackerFactory cryptocomFactory
@inject IGateIoTrackerFactory gateioFactory
@inject IHTXTrackerFactory htxFactory
@inject IKrakenTrackerFactory krakenFactory
@inject IKucoinTrackerFactory kucoinFactory
@inject IMexcTrackerFactory mexcFactory
@inject IOKXTrackerFactory okxFactory
@implements IDisposable
<h3>ETH-BTC trade Trackers, live updates:</h3>
<div style="display:flex; flex-wrap: wrap;">
@foreach (var tracker in _trackers.OrderBy(p => p.Exchange))
{
<div style="margin-bottom: 20px; flex: 1; min-width: 700px;">
<h4>@tracker.Exchange</h4>
@foreach(var line in GetInfo(tracker))
{
<div>@line</div>
}
</div>
}
</div>
@code{
private List<ITradeTracker> _trackers = new List<ITradeTracker>();
private Timer _timer;
protected override async Task OnInitializedAsync()
{
var symbol = new SharedSymbol(TradingMode.Spot, "BTC", "USDT");
_trackers = new List<ITradeTracker>
{
{ binanceFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ bingXFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ bitfinexFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ bitgetFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ bitmartFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ bybitFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ coinbaseFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ coinExFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ cryptocomFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ gateioFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ htxFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ krakenFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ kucoinFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ mexcFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
{ okxFactory.CreateTradeTracker(symbol, period: TimeSpan.FromMinutes(5)) },
};
await Task.WhenAll(_trackers.Select(b => b.StartAsync()));
// Use a manual update timer so the page isn't refreshed too often
_timer = new Timer(500);
_timer.Start();
_timer.Elapsed += (o, e) => InvokeAsync(StateHasChanged);
}
private string[] GetInfo(ITradeTracker tracker)
{
var secondLastMinute = tracker.GetStats(DateTime.UtcNow.AddMinutes(-2), DateTime.UtcNow.AddMinutes(-1));
var lastMinute = tracker.GetStats(DateTime.UtcNow.AddMinutes(-1));
var compare = lastMinute.CompareTo(secondLastMinute);
return [
$"{tracker.SymbolName} | {tracker.Status} - Synced from {tracker.SyncedFrom}",
$"Total trades: {tracker.Count}",
$"Trades last minute: {lastMinute.TradeCount}, minute before: {secondLastMinute.TradeCount}",
$"Average weighted price: {lastMinute.VolumeWeightedAveragePrice}, minute before: {secondLastMinute.VolumeWeightedAveragePrice}, dif: {compare.VolumeWeightedAveragePriceDif.PercentageDifference}%"
];
}
public void Dispose()
{
_timer.Stop();
_timer.Dispose();
foreach (var tracker in _trackers.Where(b => b.Status != CryptoExchange.Net.Objects.SyncStatus.Disconnected))
// It's not necessary to wait for this
_ = tracker.StopAsync();
}
}
+6 -1
View File
@@ -14,7 +14,7 @@
</li> </li>
<li class="nav-item px-3"> <li class="nav-item px-3">
<NavLink class="nav-link" href="SpotClient"> <NavLink class="nav-link" href="SpotClient">
Get data ISpotClient Get data SharedClient
</NavLink> </NavLink>
</li> </li>
<li class="nav-item px-3"> <li class="nav-item px-3">
@@ -27,6 +27,11 @@
Order books Order books
</NavLink> </NavLink>
</li> </li>
<li class="nav-item px-3">
<NavLink class="nav-link" href="Trackers">
Trackers
</NavLink>
</li>
</ul> </ul>
</div> </div>
+1
View File
@@ -43,6 +43,7 @@ namespace BlazorClient
services.AddBybit(); services.AddBybit();
services.AddCoinbase(); services.AddCoinbase();
services.AddCoinEx(); services.AddCoinEx();
services.AddCryptoCom();
services.AddGateIo(); services.AddGateIo();
services.AddHTX(); services.AddHTX();
services.AddKraken(); services.AddKraken();
+1
View File
@@ -16,6 +16,7 @@
@using Bybit.Net.Interfaces.Clients; @using Bybit.Net.Interfaces.Clients;
@using Coinbase.Net.Interfaces.Clients; @using Coinbase.Net.Interfaces.Clients;
@using CoinEx.Net.Interfaces.Clients; @using CoinEx.Net.Interfaces.Clients;
@using CryptoCom.Net.Interfaces.Clients;
@using GateIo.Net.Interfaces.Clients; @using GateIo.Net.Interfaces.Clients;
@using HTX.Net.Interfaces.Clients; @using HTX.Net.Interfaces.Clients;
@using Kraken.Net.Interfaces.Clients; @using Kraken.Net.Interfaces.Clients;
+14 -13
View File
@@ -6,19 +6,20 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="10.5.0" /> <PackageReference Include="Binance.Net" Version="10.8.0" />
<PackageReference Include="Bitfinex.Net" Version="7.8.0" /> <PackageReference Include="Bitfinex.Net" Version="7.9.0" />
<PackageReference Include="BitMart.Net" Version="1.2.0" /> <PackageReference Include="BitMart.Net" Version="1.5.0" />
<PackageReference Include="Bybit.Net" Version="3.14.0" /> <PackageReference Include="Bybit.Net" Version="3.15.0" />
<PackageReference Include="CoinEx.Net" Version="7.7.0" /> <PackageReference Include="CoinEx.Net" Version="7.8.0" />
<PackageReference Include="GateIo.Net" Version="1.6.0" /> <PackageReference Include="CryptoCom.Net" Version="1.1.0" />
<PackageReference Include="JK.Bitget.Net" Version="1.10.0" /> <PackageReference Include="GateIo.Net" Version="1.10.0" />
<PackageReference Include="JK.Mexc.Net" Version="1.8.0" /> <PackageReference Include="JK.Bitget.Net" Version="1.11.0" />
<PackageReference Include="JK.OKX.Net" Version="2.4.0" /> <PackageReference Include="JK.Mexc.Net" Version="1.10.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="1.0.0" /> <PackageReference Include="JK.OKX.Net" Version="2.7.0" />
<PackageReference Include="JKorf.HTX.Net" Version="6.1.0" /> <PackageReference Include="JKorf.Coinbase.Net" Version="1.2.0" />
<PackageReference Include="KrakenExchange.Net" Version="4.12.0" /> <PackageReference Include="JKorf.HTX.Net" Version="6.3.0" />
<PackageReference Include="Kucoin.Net" Version="5.14.0" /> <PackageReference Include="KrakenExchange.Net" Version="5.1.0" />
<PackageReference Include="Kucoin.Net" Version="5.17.0" />
</ItemGroup> </ItemGroup>
</Project> </Project>
@@ -27,7 +27,7 @@ namespace ConsoleClient.Exchanges
{ {
using var client = new BybitRestClient(); using var client = new BybitRestClient();
var result = await client.V5Api.Account.GetBalancesAsync(Bybit.Net.Enums.AccountType.Spot); var result = await client.V5Api.Account.GetBalancesAsync(Bybit.Net.Enums.AccountType.Spot);
return result.Data.List.First().Assets.ToDictionary(d => d.Asset, d => d.WalletBalance); return result.Data.List.First().Assets.ToDictionary(d => d.Asset, d => d.WalletBalance ?? 0);
} }
public async Task<IEnumerable<OpenOrder>> GetOpenOrders() public async Task<IEnumerable<OpenOrder>> GetOpenOrders()
+3 -3
View File
@@ -8,9 +8,9 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Binance.Net" Version="10.5.0" /> <PackageReference Include="Binance.Net" Version="10.8.0" />
<PackageReference Include="BitMart.Net" Version="1.2.0" /> <PackageReference Include="BitMart.Net" Version="1.5.0" />
<PackageReference Include="JK.OKX.Net" Version="2.4.0" /> <PackageReference Include="JK.OKX.Net" Version="2.7.0" />
</ItemGroup> </ItemGroup>
</Project> </Project>
+23
View File
@@ -21,6 +21,7 @@ The following API's are directly supported. Note that there are 3rd party implem
|Coinbase|[JKorf/Coinbase.Net](https://github.com/JKorf/Coinbase.Net)|[![Nuget version](https://img.shields.io/nuget/v/JKorf.Coinbase.Net.svg?style=flat-square)](https://www.nuget.org/packages/JKorf.Coinbase.Net)| |Coinbase|[JKorf/Coinbase.Net](https://github.com/JKorf/Coinbase.Net)|[![Nuget version](https://img.shields.io/nuget/v/JKorf.Coinbase.Net.svg?style=flat-square)](https://www.nuget.org/packages/JKorf.Coinbase.Net)|
|CoinEx|[JKorf/CoinEx.Net](https://github.com/JKorf/CoinEx.Net)|[![Nuget version](https://img.shields.io/nuget/v/CoinEx.net.svg?style=flat-square)](https://www.nuget.org/packages/CoinEx.Net)| |CoinEx|[JKorf/CoinEx.Net](https://github.com/JKorf/CoinEx.Net)|[![Nuget version](https://img.shields.io/nuget/v/CoinEx.net.svg?style=flat-square)](https://www.nuget.org/packages/CoinEx.Net)|
|CoinGecko|[JKorf/CoinGecko.Net](https://github.com/JKorf/CoinGecko.Net)|[![Nuget version](https://img.shields.io/nuget/v/CoinGecko.net.svg?style=flat-square)](https://www.nuget.org/packages/CoinGecko.Net)| |CoinGecko|[JKorf/CoinGecko.Net](https://github.com/JKorf/CoinGecko.Net)|[![Nuget version](https://img.shields.io/nuget/v/CoinGecko.net.svg?style=flat-square)](https://www.nuget.org/packages/CoinGecko.Net)|
|Crypto.com|[JKorf/CryptoCom.Net](https://github.com/JKorf/CryptoCom.Net)|[![Nuget version](https://img.shields.io/nuget/v/CryptoCom.net.svg?style=flat-square)](https://www.nuget.org/packages/CryptoCom.Net)|
|Gate.io|[JKorf/GateIo.Net](https://github.com/JKorf/GateIo.Net)|[![Nuget version](https://img.shields.io/nuget/v/GateIo.net.svg?style=flat-square)](https://www.nuget.org/packages/GateIo.Net)| |Gate.io|[JKorf/GateIo.Net](https://github.com/JKorf/GateIo.Net)|[![Nuget version](https://img.shields.io/nuget/v/GateIo.net.svg?style=flat-square)](https://www.nuget.org/packages/GateIo.Net)|
|HTX|[JKorf/HTX.Net](https://github.com/JKorf/HTX.Net)|[![Nuget version](https://img.shields.io/nuget/v/JKorf.HTX.net.svg?style=flat-square)](https://www.nuget.org/packages/JKorf.HTX.Net)| |HTX|[JKorf/HTX.Net](https://github.com/JKorf/HTX.Net)|[![Nuget version](https://img.shields.io/nuget/v/JKorf.HTX.net.svg?style=flat-square)](https://www.nuget.org/packages/JKorf.HTX.Net)|
|Kraken|[JKorf/Kraken.Net](https://github.com/JKorf/Kraken.Net)|[![Nuget version](https://img.shields.io/nuget/v/KrakenExchange.net.svg?style=flat-square)](https://www.nuget.org/packages/KrakenExchange.Net)| |Kraken|[JKorf/Kraken.Net](https://github.com/JKorf/Kraken.Net)|[![Nuget version](https://img.shields.io/nuget/v/KrakenExchange.net.svg?style=flat-square)](https://www.nuget.org/packages/KrakenExchange.Net)|
@@ -48,6 +49,28 @@ 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 8.1.1 - 01 Nov 2024
* Fixed socket connections trying to authenticated connection when it's marked as dedicated request connection even when no authentication is needed
* Fixed System.Text.Json ArrayConverter not passing serializer options to nested deserialization
* Fixed System.Text.Json ArrayConverter creating new serializer options each time a JsonConverter attribute is encountered
* Version 8.1.0 - 28 Oct 2024
* Added KlineTracker and TradeTracker implementation
* Added Side to SharedTrade model
* Added overload for Create method in OrderBookFactory using SharedSymbol
* Added ValidateMessage method to websocket Query object to filter messages even though it is matched to the query based on the ListenIdentifier
* Added DoHandleReset method for websocket subscriptions
* Added ConnectionId to RequestDefinition to correctly handle connection and path rate limiting configuration
* Added System.Text.Json ArrayConverter Write implementation
* Updated SharedFuturesTicker LastPrice, HighPrice and LowPrice properties to be nullable
* Updated SetApiCredentials method to also updated the credentials on the client specific options to prevent unknown client credentials in some situations
* Version 8.0.3 - 14 Oct 2024
* Added support for duplicate array indexes in System.Text.Json ArrayConverter
* Added fallback for unparsable value in System.Text.Json NumberStringConverter
* Added Authenticated property on base client and shared client
* Added GetValues System.Text.Json implementation in message accessor
* Version 8.0.2 - 09 Oct 2024 * Version 8.0.2 - 09 Oct 2024
* Updated dependency versions, including System.Text.Json from 8.0.4 to 8.0.5 containing a vulnerability fix * Updated dependency versions, including System.Text.Json from 8.0.4 to 8.0.5 containing a vulnerability fix
+835 -6
View File
File diff suppressed because it is too large Load Diff