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

Compare commits

...

32 Commits

Author SHA1 Message Date
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
Jkorf 84d36544e4 Updated to version 8.0.2 2024-10-09 08:49:19 +02:00
Jkorf a71f57ae7f Updated dependency versions, including System.Text.Json containing a vulnerability fix 2024-10-09 08:46:15 +02:00
Jkorf 6e5bcd5e9a Some small doc fixes 2024-10-07 15:30:09 +02:00
Jkorf 4131e563c3 Added Coinbase reference, updated examples 2024-10-07 15:20:12 +02:00
Jkorf 613766dbca Updated to version 8.0.1 2024-10-07 13:03:08 +02:00
Jkorf 23b07d709e Added testing check for year 1 datetime 2024-10-07 12:57:23 +02:00
Jkorf bbbdac2fd3 Added ToRfc3339String extension method for DateTime 2024-10-07 12:57:04 +02:00
Jkorf c614b7869c Added check for datetime year 1 to be deserialized as null 2024-10-04 18:50:27 +02:00
Jkorf 1f31e4a9d7 Added cached lib versions properties 2024-10-04 18:49:58 +02:00
Jkorf 6cb6cd6b11 Note 2024-10-03 16:33:00 +02:00
Jkorf 17ffec329f Fixed typo 2024-09-27 15:15:09 +02:00
Jkorf 7a3927ef49 Add parameters documentation for shared clients 2024-09-27 15:13:50 +02:00
Jkorf c1b0437c93 Updated examples 2024-09-27 13:53:33 +02:00
51 changed files with 3319 additions and 123 deletions
@@ -302,7 +302,7 @@ namespace CryptoExchange.Net.UnitTests
public async Task ApiKeyRateLimiterBasics(string key1, string key2, string endpoint1, string endpoint2, bool expectLimited)
{
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 requestDefinition2 = new RequestDefinition(endpoint2, HttpMethod.Get) { Authenticated = key2 != null };
@@ -5,6 +5,8 @@ using NUnit.Framework;
using System;
using System.Text.Json.Serialization;
using NUnit.Framework.Legacy;
using CryptoExchange.Net.Converters;
using CryptoExchange.Net.Testing.Comparers;
namespace CryptoExchange.Net.UnitTests
{
@@ -242,6 +244,44 @@ namespace CryptoExchange.Net.UnitTests
var result = JsonSerializer.Deserialize<STJDecimalObject>("{ \"test\": " + value + "}");
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
@@ -281,4 +321,42 @@ namespace CryptoExchange.Net.UnitTests
[JsonConverter(typeof(BoolConverter))]
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; }
}
}
+8 -1
View File
@@ -11,7 +11,9 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "BlazorClient", "Examples\Bl
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Examples", "Examples", "{5734C2A9-F12C-4754-A8B9-640C24DC4E02}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ConsoleClient", "Examples\ConsoleClient\ConsoleClient.csproj", "{23480C58-23BF-4EBF-A173-B7F51A043A99}"
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "ConsoleClient", "Examples\ConsoleClient\ConsoleClient.csproj", "{23480C58-23BF-4EBF-A173-B7F51A043A99}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "SharedClients", "Examples\SharedClients\SharedClients.csproj", "{988A87EF-EAEA-4313-A6CF-FA869813D5AB}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
@@ -35,6 +37,10 @@ Global
{23480C58-23BF-4EBF-A173-B7F51A043A99}.Debug|Any CPU.Build.0 = Debug|Any CPU
{23480C58-23BF-4EBF-A173-B7F51A043A99}.Release|Any CPU.ActiveCfg = Release|Any CPU
{23480C58-23BF-4EBF-A173-B7F51A043A99}.Release|Any CPU.Build.0 = Release|Any CPU
{988A87EF-EAEA-4313-A6CF-FA869813D5AB}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{988A87EF-EAEA-4313-A6CF-FA869813D5AB}.Debug|Any CPU.Build.0 = Debug|Any CPU
{988A87EF-EAEA-4313-A6CF-FA869813D5AB}.Release|Any CPU.ActiveCfg = Release|Any CPU
{988A87EF-EAEA-4313-A6CF-FA869813D5AB}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@@ -42,6 +48,7 @@ Global
GlobalSection(NestedProjects) = preSolution
{AF4F5C19-162E-48F4-8B0B-BA5A2D7CE06A} = {5734C2A9-F12C-4754-A8B9-640C24DC4E02}
{23480C58-23BF-4EBF-A173-B7F51A043A99} = {5734C2A9-F12C-4754-A8B9-640C24DC4E02}
{988A87EF-EAEA-4313-A6CF-FA869813D5AB} = {5734C2A9-F12C-4754-A8B9-640C24DC4E02}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {0D1B9CE9-E0B7-4B8B-88BF-6EA2CC8CA3D7}
@@ -11,12 +11,12 @@ namespace CryptoExchange.Net.Authentication
public class ApiCredentials
{
/// <summary>
/// The api key to authenticate requests
/// The api key / label to authenticate requests
/// </summary>
public string Key { get; }
/// <summary>
/// The api secret to authenticate requests
/// The api secret or private key to authenticate requests
/// </summary>
public string Secret { get; }
@@ -28,8 +28,8 @@ namespace CryptoExchange.Net.Authentication
/// <summary>
/// Create Api credentials providing an api key and secret for authentication
/// </summary>
/// <param name="key">The api key used for identification</param>
/// <param name="secret">The api secret used for signing</param>
/// <param name="key">The api key / label used for identification</param>
/// <param name="secret">The api secret or private key used for signing</param>
public ApiCredentials(string key, string secret) : this(key, secret, ApiCredentialsType.Hmac)
{
}
@@ -37,8 +37,8 @@ namespace CryptoExchange.Net.Authentication
/// <summary>
/// Create Api credentials providing an api key and secret for authentication
/// </summary>
/// <param name="key">The api key used for identification</param>
/// <param name="secret">The api secret used for signing</param>
/// <param name="key">The api key / label used for identification</param>
/// <param name="secret">The api secret or private key used for signing</param>
/// <param name="credentialsType">The type of credentials</param>
public ApiCredentials(string key, string secret, ApiCredentialsType credentialsType)
{
@@ -38,6 +38,11 @@ namespace CryptoExchange.Net.Clients
/// </summary>
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>
/// Api options
/// </summary>
@@ -83,6 +88,7 @@ namespace CryptoExchange.Net.Clients
/// <inheritdoc />
public void SetApiCredentials<T>(T credentials) where T : ApiCredentials
{
ApiOptions.ApiCredentials = credentials;
if (credentials != null)
AuthenticationProvider = CreateAuthenticationProvider(credentials.Copy());
}
+26 -1
View File
@@ -12,6 +12,28 @@ namespace CryptoExchange.Net.Clients
/// </summary>
public abstract class BaseClient : IDisposable
{
/// <summary>
/// Version of the CryptoExchange.Net base library
/// </summary>
public Version CryptoExchangeLibVersion { get; } = typeof(BaseClient).Assembly.GetName().Version;
/// <summary>
/// Version of the client implementation
/// </summary>
public Version ExchangeLibVersion
{
get
{
lock(_versionLock)
{
if (_exchangeVersion == null)
_exchangeVersion = GetType().Assembly.GetName().Version;
return _exchangeVersion;
}
}
}
/// <summary>
/// The name of the API the client is for
/// </summary>
@@ -27,6 +49,9 @@ namespace CryptoExchange.Net.Clients
/// </summary>
protected internal ILogger _logger;
private object _versionLock = new object();
private Version _exchangeVersion;
/// <summary>
/// Provided client options
/// </summary>
@@ -57,7 +82,7 @@ namespace CryptoExchange.Net.Clients
throw new ArgumentNullException(nameof(options));
ClientOptions = options;
_logger.Log(LogLevel.Trace, $"Client configuration: {options}, CryptoExchange.Net: v{typeof(BaseClient).Assembly.GetName().Version}, {Exchange}.Net: v{GetType().Assembly.GetName().Version}");
_logger.Log(LogLevel.Trace, $"Client configuration: {options}, CryptoExchange.Net: v{CryptoExchangeLibVersion}, {Exchange}.Net: v{ExchangeLibVersion}");
}
/// <summary>
@@ -42,8 +42,63 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
public override void Write(Utf8JsonWriter writer, T value, JsonSerializerOptions options)
{
// TODO
throw new NotImplementedException();
if (value == null)
{
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 />
@@ -56,6 +111,19 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
return (T)ParseObject(ref reader, result, typeToConvert);
}
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)
{
var attributes = new List<ArrayPropertyInfo>();
@@ -71,7 +139,7 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
ArrayProperty = att,
PropertyInfo = property,
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
});
}
@@ -94,42 +162,46 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
if (reader.TokenType == JsonTokenType.EndArray)
break;
var attribute = attributes.SingleOrDefault(a => a.ArrayProperty.Index == index);
if (attribute == null)
var indexAttributes = attributes.Where(a => a.ArrayProperty.Index == index);
if (!indexAttributes.Any())
{
index++;
continue;
}
var targetType = attribute.TargetType;
object? value = null;
if (attribute.JsonConverterType != null)
foreach (var attribute in indexAttributes)
{
// Has JsonConverter attribute
var options = new JsonSerializerOptions();
options.Converters.Add((JsonConverter)Activator.CreateInstance(attribute.JsonConverterType));
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
var targetType = attribute.TargetType;
object? value = null;
if (attribute.JsonConverterType != null)
{
JsonTokenType.Null => null,
JsonTokenType.False => false,
JsonTokenType.True => true,
JsonTokenType.String => reader.GetString(),
JsonTokenType.Number => reader.GetDecimal(),
_ => throw new NotImplementedException($"Array deserialization of type {reader.TokenType} not supported"),
};
// Has JsonConverter attribute
var options = new JsonSerializerOptions();
options.Converters.Add((JsonConverter)Activator.CreateInstance(attribute.JsonConverterType));
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,
JsonTokenType.False => false,
JsonTokenType.True => true,
JsonTokenType.String => reader.GetString(),
JsonTokenType.Number => reader.GetDecimal(),
JsonTokenType.StartObject => JsonSerializer.Deserialize(ref reader, attribute.TargetType),
_ => throw new NotImplementedException($"Array deserialization of type {reader.TokenType} not supported"),
};
}
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++;
}
@@ -57,6 +57,7 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
var stringValue = reader.GetString();
if (string.IsNullOrWhiteSpace(stringValue)
|| stringValue == "-1"
|| stringValue == "0001-01-01T00:00:00Z"
|| double.TryParse(stringValue, out var doubleVal) && doubleVal == 0)
{
return default;
@@ -23,7 +23,14 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
return reader.GetDecimal().ToString();
}
return reader.GetString();
try
{
return reader.GetString();
}
catch (Exception)
{
return null;
}
}
/// <inheritdoc />
@@ -133,7 +133,20 @@ namespace CryptoExchange.Net.Converters.SystemTextJson
}
/// <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)
{
+7 -7
View File
@@ -6,9 +6,9 @@
<PackageId>CryptoExchange.Net</PackageId>
<Authors>JKorf</Authors>
<Description>CryptoExchange.Net is a base library which is used to implement different cryptocurrency (exchange) API's. It provides a standardized way of implementing different API's, which results in a very similar experience for users of the API implementations.</Description>
<PackageVersion>8.0.0</PackageVersion>
<AssemblyVersion>8.0.0</AssemblyVersion>
<FileVersion>8.0.0</FileVersion>
<PackageVersion>8.1.0</PackageVersion>
<AssemblyVersion>8.1.0</AssemblyVersion>
<FileVersion>8.1.0</FileVersion>
<PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance>
<PackageTags>OKX;OKX.Net;Mexc;Mexc.Net;Kucoin;Kucoin.Net;Kraken;Kraken.Net;Huobi;Huobi.Net;CoinEx;CoinEx.Net;Bybit;Bybit.Net;Bitget;Bitget.Net;Bitfinex;Bitfinex.Net;Binance;Binance.Net;CryptoCurrency;CryptoCurrency Exchange</PackageTags>
<RepositoryType>git</RepositoryType>
@@ -52,12 +52,12 @@
<PrivateAssets>all</PrivateAssets>
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
<PackageReference Include="Microsoft.Extensions.Http" Version="8.0.0" />
<PackageReference Include="Microsoft.Extensions.Http" Version="8.0.1" />
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="8.0.1" />
<PackageReference Include="Microsoft.Extensions.Logging.Abstractions" Version="8.0.1" />
<PackageReference Include="System.Text.Json" Version="8.0.4" />
<PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="8.0.2" />
<PackageReference Include="Microsoft.Extensions.Logging.Abstractions" Version="8.0.2" />
<PackageReference Include="System.Text.Json" Version="8.0.5" />
</ItemGroup>
</Project>
+10
View File
@@ -189,6 +189,16 @@ namespace CryptoExchange.Net
throw new ArgumentException($"No values provided for parameter {argumentName}", argumentName);
}
/// <summary>
/// Format a string to RFC3339/ISO8601 string
/// </summary>
/// <param name="dateTime"></param>
/// <returns></returns>
public static string ToRfc3339String(this DateTime dateTime)
{
return dateTime.ToString("yyyy-MM-dd'T'HH:mm:ss.fffzzz", DateTimeFormatInfo.InvariantInfo);
}
/// <summary>
/// Format an exception and inner exception to a readable string
/// </summary>
@@ -1,4 +1,5 @@
using CryptoExchange.Net.Objects.Options;
using CryptoExchange.Net.SharedApis;
using System;
namespace CryptoExchange.Net.Interfaces
@@ -23,5 +24,12 @@ namespace CryptoExchange.Net.Interfaces
/// <param name="options">Options for the order book</param>
/// <returns></returns>
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);
/// <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>
/// <returns></returns>
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
}
/// <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>
/// Status of the order book
/// </summary>
@@ -58,12 +58,16 @@ namespace CryptoExchange.Net.Objects
/// </summary>
public IRateLimitGuard? LimitGuard { get; set; }
/// <summary>
/// Whether this request should never be cached
/// </summary>
public bool PreventCaching { get; set; }
/// <summary>
/// Connection id
/// </summary>
public int? ConnectionId { get; set; }
/// <summary>
/// ctor
/// </summary>
@@ -1,5 +1,6 @@
using CryptoExchange.Net.Interfaces;
using CryptoExchange.Net.Objects.Options;
using CryptoExchange.Net.SharedApis;
using System;
namespace CryptoExchange.Net.OrderBook
@@ -8,14 +9,14 @@ namespace CryptoExchange.Net.OrderBook
public class OrderBookFactory<TOptions> : IOrderBookFactory<TOptions> where TOptions: OrderBookOptions
{
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>
/// ctor
/// </summary>
/// <param name="symbolCtor"></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;
_assetsCtor = assetsCtor;
@@ -25,6 +26,9 @@ namespace CryptoExchange.Net.OrderBook
public ISymbolOrderBook Create(string symbol, Action<TOptions>? options = null) => _symbolCtor(symbol, options);
/// <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>
public static Func<RequestDefinition, string, string?, string> PerEndpoint { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => def.Path + def.Method);
/// <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
/// </summary>
public static Func<RequestDefinition, string, string?, string> PerApiKey { get; } = new Func<RequestDefinition, string, string?, string>((def, host, key) => key!);
@@ -18,6 +18,11 @@ namespace CryptoExchange.Net.SharedApis
/// </summary>
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>
/// Format a base and quote asset to an exchange accepted symbol
/// </summary>
@@ -14,15 +14,15 @@ namespace CryptoExchange.Net.SharedApis
/// <summary>
/// Last trade price
/// </summary>
public decimal LastPrice { get; set; }
public decimal? LastPrice { get; set; }
/// <summary>
/// High price in the last 24h
/// </summary>
public decimal HighPrice { get; set; }
public decimal? HighPrice { get; set; }
/// <summary>
/// Low price in the last 24h
/// </summary>
public decimal LowPrice { get; set; }
public decimal? LowPrice { get; set; }
/// <summary>
/// The volume in the last 24h
/// </summary>
@@ -51,7 +51,7 @@ namespace CryptoExchange.Net.SharedApis
/// <summary>
/// ctor
/// </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;
LastPrice = lastPrice;
@@ -19,6 +19,10 @@ namespace CryptoExchange.Net.SharedApis
/// Trade time
/// </summary>
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>
/// ctor
@@ -209,7 +209,7 @@ namespace CryptoExchange.Net.Sockets
{
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);
if (!limitResult)
return new CallResult(new ClientRateLimitError("Connection limit reached"));
@@ -475,7 +475,7 @@ namespace CryptoExchange.Net.Sockets
/// <returns></returns>
private async Task SendLoopAsync()
{
var requestDefinition = new RequestDefinition(Id.ToString(), HttpMethod.Get);
var requestDefinition = new RequestDefinition(Uri.AbsolutePath, HttpMethod.Get) { ConnectionId = Id };
try
{
while (true)
+12 -1
View File
@@ -177,6 +177,10 @@ namespace CryptoExchange.Net.Sockets
/// <inheritdoc />
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++;
if (CurrentResponses == RequiredResponses)
{
@@ -186,7 +190,7 @@ namespace CryptoExchange.Net.Sockets
if (Result?.Success != false)
// 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)
{
@@ -198,6 +202,13 @@ namespace CryptoExchange.Net.Sockets
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>
/// Handle the query response
/// </summary>
@@ -268,7 +268,7 @@ namespace CryptoExchange.Net.Sockets
lock (_listenersLock)
{
foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription))
subscription.Confirmed = false;
subscription.Reset();
foreach (var query in _listeners.OfType<Query>().ToList())
{
@@ -293,7 +293,7 @@ namespace CryptoExchange.Net.Sockets
lock (_listenersLock)
{
foreach (var subscription in _listeners.OfType<Subscription>().Where(l => l.UserSubscription))
subscription.Confirmed = false;
subscription.Reset();
foreach (var query in _listeners.OfType<Query>().ToList())
{
@@ -506,8 +506,11 @@ namespace CryptoExchange.Net.Sockets
}
if (processor is Subscription subscriptionProcessor && !subscriptionProcessor.Confirmed)
{
// If this message is for this listener then it is automatically confirmed, even if the subscription is not (yet) confirmed
subscriptionProcessor.Confirmed = true;
// This doesn't trigger a waiting subscribe query, should probably also somehow set the wait event for that
}
// 6. Deserialize the message
object? deserialized = null;
@@ -130,6 +130,20 @@ namespace CryptoExchange.Net.Sockets
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>
/// Handle the update message
/// </summary>
@@ -188,7 +188,7 @@ namespace CryptoExchange.Net.Testing.Comparers
{
if (propertyValue == default && propValue.Type != JTokenType.Null && !string.IsNullOrEmpty(propValue.ToString()))
{
if (propertyType == typeof(DateTime?) && (propValue.ToString() == "" || propValue.ToString() == "0" || propValue.ToString() == "-1"))
if (propertyType == typeof(DateTime?) && (propValue.ToString() == "" || propValue.ToString() == "0" || propValue.ToString() == "-1" || propValue.ToString() == "01/01/0001 00:00:00"))
return;
// Property value not correct
@@ -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 -11
View File
@@ -5,17 +5,22 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="9.7.1" />
<PackageReference Include="Bitfinex.Net" Version="7.2.2" />
<PackageReference Include="Bybit.Net" Version="3.7.1" />
<PackageReference Include="CoinEx.Net" Version="6.2.1" />
<PackageReference Include="Huobi.Net" Version="5.2.1" />
<PackageReference Include="JK.BingX.Net" Version="1.0.0" />
<PackageReference Include="JK.Bitget.Net" Version="1.3.1" />
<PackageReference Include="JK.OKX.Net" Version="1.7.1" />
<PackageReference Include="KrakenExchange.Net" Version="4.4.3" />
<PackageReference Include="Kucoin.Net" Version="5.3.2" />
<PackageReference Include="Serilog.AspNetCore" Version="6.0.0" />
<PackageReference Include="Binance.Net" Version="10.7.0" />
<PackageReference Include="Bitfinex.Net" Version="7.8.2" />
<PackageReference Include="BitMart.Net" Version="1.4.0" />
<PackageReference Include="Bybit.Net" Version="3.14.3" />
<PackageReference Include="CoinEx.Net" Version="7.7.2" />
<PackageReference Include="CryptoCom.Net" Version="1.0.1" />
<PackageReference Include="GateIo.Net" Version="1.9.0" />
<PackageReference Include="JK.BingX.Net" Version="1.11.2" />
<PackageReference Include="JK.Bitget.Net" Version="1.10.4" />
<PackageReference Include="JK.Mexc.Net" Version="1.9.0" />
<PackageReference Include="JK.OKX.Net" Version="2.6.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="1.1.2" />
<PackageReference Include="JKorf.HTX.Net" Version="6.2.0" />
<PackageReference Include="KrakenExchange.Net" Version="5.0.2" />
<PackageReference Include="Kucoin.Net" Version="5.16.0" />
<PackageReference Include="Serilog.AspNetCore" Version="8.0.2" />
</ItemGroup>
</Project>
+30 -5
View File
@@ -2,12 +2,17 @@
@inject IBinanceRestClient binanceClient
@inject IBingXRestClient bingXClient
@inject IBitfinexRestClient bitfinexClient
@inject IBitMartRestClient bitmartClient
@inject IBitgetRestClient bitgetClient
@inject IBybitRestClient bybitClient
@inject ICoinbaseRestClient coinbaseClient
@inject ICoinExRestClient coinexClient
@inject IHuobiRestClient huobiClient
@inject ICryptoComRestClient cryptocomClient
@inject IGateIoRestClient gateioClient
@inject IHTXRestClient huobiClient
@inject IKrakenRestClient krakenClient
@inject IKucoinRestClient kucoinClient
@inject IMexcRestClient mexcClient
@inject IOKXRestClient okxClient
<h3>BTC-USD prices:</h3>
@@ -25,14 +30,19 @@
var bingXTask = bingXClient.SpotApi.ExchangeData.GetTickersAsync("BTC-USDT");
var bitfinexTask = bitfinexClient.SpotApi.ExchangeData.GetTickerAsync("tBTCUSD");
var bitgetTask = bitgetClient.SpotApi.ExchangeData.GetTickerAsync("BTCUSDT_SPBL");
var bitmartTask = bitmartClient.SpotApi.ExchangeData.GetTickerAsync("BTC_USDT");
var bybitTask = bybitClient.V5Api.ExchangeData.GetSpotTickersAsync("BTCUSDT");
var coinbaseTask = coinbaseClient.AdvancedTradeApi.ExchangeData.GetSymbolAsync("BTC-USDT");
var coinexTask = coinexClient.SpotApi.ExchangeData.GetTickerAsync("BTCUSDT");
var huobiTask = huobiClient.SpotApi.ExchangeData.GetTickerAsync("btcusdt");
var cryptocomTask = cryptocomClient.ExchangeApi.ExchangeData.GetTickersAsync("BTC_USDT");
var gateioTask = gateioClient.SpotApi.ExchangeData.GetTickersAsync("BTC_USDT");
var htxTask = huobiClient.SpotApi.ExchangeData.GetTickerAsync("btcusdt");
var krakenTask = krakenClient.SpotApi.ExchangeData.GetTickerAsync("XBTUSD");
var kucoinTask = kucoinClient.SpotApi.ExchangeData.GetTickerAsync("BTC-USDT");
var mexcTask = mexcClient.SpotApi.ExchangeData.GetTickerAsync("BTCUSDT");
var okxTask = okxClient.UnifiedApi.ExchangeData.GetTickerAsync("BTCUSDT");
await Task.WhenAll(binanceTask, bingXTask, bitfinexTask, bybitTask, coinexTask, huobiTask, krakenTask, kucoinTask);
await Task.WhenAll(binanceTask, bingXTask, bitfinexTask, bitgetTask, bitmartTask, bybitTask, coinexTask, gateioTask, htxTask, krakenTask, kucoinTask, mexcTask, okxTask);
if (binanceTask.Result.Success)
_prices.Add("Binance", binanceTask.Result.Data.LastPrice);
@@ -46,14 +56,26 @@
if (bitgetTask.Result.Success)
_prices.Add("Bitget", bitgetTask.Result.Data.ClosePrice);
if (bitmartTask.Result.Success)
_prices.Add("BitMart", bitgetTask.Result.Data.ClosePrice);
if (bybitTask.Result.Success)
_prices.Add("Bybit", bybitTask.Result.Data.List.First().LastPrice);
if (coinbaseTask.Result.Success)
_prices.Add("Coinbase", coinbaseTask.Result.Data.LastPrice ?? 0);
if (coinexTask.Result.Success)
_prices.Add("CoinEx", coinexTask.Result.Data.Ticker.LastPrice);
if (huobiTask.Result.Success)
_prices.Add("Huobi", huobiTask.Result.Data.ClosePrice ?? 0);
if (cryptocomTask.Result.Success)
_prices.Add("CryptoCom", cryptocomTask.Result.Data.First().LastPrice ?? 0);
if (gateioTask.Result.Success)
_prices.Add("GateIo", gateioTask.Result.Data.First().LastPrice);
if (htxTask.Result.Success)
_prices.Add("HTX", htxTask.Result.Data.ClosePrice ?? 0);
if (krakenTask.Result.Success)
_prices.Add("Kraken", krakenTask.Result.Data.First().Value.LastTrade.Price);
@@ -61,6 +83,9 @@
if (kucoinTask.Result.Success)
_prices.Add("Kucoin", kucoinTask.Result.Data.LastPrice ?? 0);
if (mexcTask.Result.Success)
_prices.Add("Mexc", mexcTask.Result.Data.LastPrice);
if (okxTask.Result.Success)
_prices.Add("OKX", okxTask.Result.Data.LastPrice ?? 0);
}
+13 -3
View File
@@ -3,11 +3,16 @@
@inject IBingXSocketClient bingXSocketClient
@inject IBitfinexSocketClient bitfinexSocketClient
@inject IBitgetSocketClient bitgetSocketClient
@inject IBitMartSocketClient bitmartSocketClient
@inject IBybitSocketClient bybitSocketClient
@inject ICoinbaseSocketClient coinbaseSocketClient
@inject ICoinExSocketClient coinExSocketClient
@inject IHuobiSocketClient huobiSocketClient
@inject ICryptoComSocketClient cryptocomSocketClient
@inject IGateIoSocketClient gateioSocketClient
@inject IHTXSocketClient htxSocketClient
@inject IKrakenSocketClient krakenSocketClient
@inject IKucoinSocketClient kucoinSocketClient
@inject IMexcSocketClient mexcSocketClient
@inject IOKXSocketClient okxSocketClient
@using System.Collections.Concurrent
@using CryptoExchange.Net.Objects
@@ -33,11 +38,16 @@
bingXSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("BingX", data.Data.LastPrice)),
bitfinexSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("tETHBTC", data => UpdateData("Bitfinex", data.Data.LastPrice)),
bitgetSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bitget", data.Data.LastPrice)),
bitmartSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH_BTC", data => UpdateData("BitMart", data.Data.LastPrice)),
bybitSocketClient.V5SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("Bybit", data.Data.LastPrice)),
coinExSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETHBTC", data => UpdateData("CoinEx", data.Data.LastPrice)),
huobiSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ethbtc", data => UpdateData("Huobi", data.Data.ClosePrice ?? 0)),
krakenSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH/XBT", data => UpdateData("Kraken", data.Data.LastTrade.Price)),
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)),
htxSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ethbtc", data => UpdateData("HTX", data.Data.ClosePrice ?? 0)),
krakenSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH/XBT", data => UpdateData("Kraken", data.Data.LastPrice)),
kucoinSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Kucoin", data.Data.LastPrice ?? 0)),
mexcSocketClient.SpotApi.SubscribeToMiniTickerUpdatesAsync("ETHBTC", data => UpdateData("Mexc", data.Data.LastPrice)),
okxSocketClient.UnifiedApi.ExchangeData.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("OKX", data.Data.LastPrice ?? 0)),
};
+18 -3
View File
@@ -5,23 +5,33 @@
@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 Huobi.Net.Interfaces
@using CryptoCom.Net.Interfaces
@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 IBinanceOrderBookFactory binanceFactory
@inject IBingXOrderBookFactory bingXFactory
@inject IBitfinexOrderBookFactory bitfinexFactory
@inject IBitgetOrderBookFactory bitgetFactory
@inject IBitMartOrderBookFactory bitmartFactory
@inject IBybitOrderBookFactory bybitFactory
@inject ICoinbaseOrderBookFactory coinbaseFactory
@inject ICoinExOrderBookFactory coinExFactory
@inject IHuobiOrderBookFactory huobiFactory
@inject ICryptoComOrderBookFactory cryptocomFactory
@inject IGateIoOrderBookFactory gateioFactory
@inject IHTXOrderBookFactory htxFactory
@inject IKrakenOrderBookFactory krakenFactory
@inject IKucoinOrderBookFactory kucoinFactory
@inject IMexcOrderBookFactory mexcFactory
@inject IOKXOrderBookFactory okxFactory
@implements IDisposable
@@ -60,11 +70,16 @@
{ "BingX", bingXFactory.CreateSpot("ETH-BTC") },
{ "Bitfinex", bitfinexFactory.Create("tETHBTC") },
{ "Bitget", bitgetFactory.CreateSpot("ETHBTC") },
{ "BitMart", bitmartFactory.CreateSpot("ETH_BTC", null) },
{ "Bybit", bybitFactory.Create("ETHBTC", Bybit.Net.Enums.Category.Spot) },
{ "Coinbase", coinbaseFactory.Create("ETH-BTC", null) },
{ "CoinEx", coinExFactory.CreateSpot("ETHBTC") },
{ "Huobi", huobiFactory.CreateSpot("ethbtc") },
{ "CryptoCom", cryptocomFactory.CreateExchange("ETH_BTC") },
{ "GateIo", gateioFactory.CreateSpot("ETH_BTC") },
{ "HTX", htxFactory.CreateSpot("ethbtc") },
{ "Kraken", krakenFactory.CreateSpot("ETH/XBT") },
{ "Kucoin", kucoinFactory.CreateSpot("ETH-BTC") },
{ "Mexc", mexcFactory.CreateSpot("ETHBTC") },
{ "OKX", okxFactory.Create("ETH-BTC") },
};
+8 -7
View File
@@ -1,5 +1,6 @@
@page "/SpotClient"
@inject ICryptoRestClient restClient
@using CryptoExchange.Net.SharedApis
@inject IEnumerable<ISpotTickerRestClient> restClients
<h3>ETH-BTC prices:</h3>
@foreach(var price in _prices.OrderBy(p => p.Key))
@@ -12,13 +13,13 @@
protected override async Task OnInitializedAsync()
{
var clients = restClient.GetSpotClients();
var tasks = clients.Select(c => (c.ExchangeName, c.GetTickerAsync(c.GetSymbolName("ETH", "BTC"))));
await Task.WhenAll(tasks.Select(t => t.Item2));
foreach(var task in tasks)
var symbol = new SharedSymbol(TradingMode.Spot, "ETH", "BTC");
var tasks = restClients.Select(x => x.GetSpotTickerAsync(new GetTickerRequest(symbol)));
await Task.WhenAll(tasks);
foreach (var ticker in tasks.Select(x => x.Result))
{
if(task.Item2.Result.Success)
_prices.Add(task.Item1, task.Item2.Result.Data.HighPrice);
if (ticker.Success)
_prices.Add(ticker.Exchange, ticker.Data.LastPrice);
}
}
+1 -1
View File
@@ -14,7 +14,7 @@
</li>
<li class="nav-item px-3">
<NavLink class="nav-link" href="SpotClient">
Get data ISpotClient
Get data SharedClient
</NavLink>
</li>
<li class="nav-item px-3">
+6 -1
View File
@@ -39,11 +39,16 @@ namespace BlazorClient
services.AddBingX();
services.AddBitfinex();
services.AddBitget();
services.AddBitMart();
services.AddBybit();
services.AddCoinbase();
services.AddCoinEx();
services.AddHuobi();
services.AddCryptoCom();
services.AddGateIo();
services.AddHTX();
services.AddKraken();
services.AddKucoin();
services.AddMexc();
services.AddOKX();
}
+6 -1
View File
@@ -12,10 +12,15 @@
@using BingX.Net.Interfaces.Clients;
@using Bitfinex.Net.Interfaces.Clients;
@using Bitget.Net.Interfaces.Clients;
@using BitMart.Net.Interfaces.Clients;
@using Bybit.Net.Interfaces.Clients;
@using Coinbase.Net.Interfaces.Clients;
@using CoinEx.Net.Interfaces.Clients;
@using Huobi.Net.Interfaces.Clients;
@using CryptoCom.Net.Interfaces.Clients;
@using GateIo.Net.Interfaces.Clients;
@using HTX.Net.Interfaces.Clients;
@using Kraken.Net.Interfaces.Clients;
@using Kucoin.Net.Interfaces.Clients;
@using Mexc.Net.Interfaces.Clients;
@using OKX.Net.Interfaces.Clients;
@using CryptoExchange.Net.Interfaces;
+15 -11
View File
@@ -1,4 +1,4 @@
<Project Sdk="Microsoft.NET.Sdk">
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
@@ -6,16 +6,20 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="9.5.0" />
<PackageReference Include="Bitfinex.Net" Version="7.1.0" />
<PackageReference Include="Bittrex.Net" Version="8.0.3" />
<PackageReference Include="Bybit.Net" Version="3.4.0" />
<PackageReference Include="CoinEx.Net" Version="6.1.0" />
<PackageReference Include="Huobi.Net" Version="5.1.0" />
<PackageReference Include="JK.Bitget.Net" Version="1.1.0" />
<PackageReference Include="JK.OKX.Net" Version="1.6.0" />
<PackageReference Include="KrakenExchange.Net" Version="4.3.0" />
<PackageReference Include="Kucoin.Net" Version="5.2.0" />
<PackageReference Include="Binance.Net" Version="10.7.0" />
<PackageReference Include="Bitfinex.Net" Version="7.8.2" />
<PackageReference Include="BitMart.Net" Version="1.4.0" />
<PackageReference Include="Bybit.Net" Version="3.14.3" />
<PackageReference Include="CoinEx.Net" Version="7.7.2" />
<PackageReference Include="CryptoCom.Net" Version="1.0.1" />
<PackageReference Include="GateIo.Net" Version="1.9.0" />
<PackageReference Include="JK.Bitget.Net" Version="1.10.4" />
<PackageReference Include="JK.Mexc.Net" Version="1.9.0" />
<PackageReference Include="JK.OKX.Net" Version="2.6.0" />
<PackageReference Include="JKorf.Coinbase.Net" Version="1.1.2" />
<PackageReference Include="JKorf.HTX.Net" Version="6.2.0" />
<PackageReference Include="KrakenExchange.Net" Version="5.0.2" />
<PackageReference Include="Kucoin.Net" Version="5.16.0" />
</ItemGroup>
</Project>
@@ -27,7 +27,7 @@ namespace ConsoleClient.Exchanges
{
using var client = new BybitRestClient();
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()
-2
View File
@@ -4,12 +4,10 @@ using System.Globalization;
using System.Linq;
using System.Threading.Tasks;
using Binance.Net.Clients;
using Binance.Net.Objects;
using Bybit.Net.Clients;
using ConsoleClient.Exchanges;
using CryptoExchange.Net.Authentication;
using CryptoExchange.Net.Objects.Sockets;
using CryptoExchange.Net.Sockets;
namespace ConsoleClient
{
+50
View File
@@ -0,0 +1,50 @@
using Binance.Net.Clients;
using BitMart.Net.Clients;
using CryptoExchange.Net.SharedApis;
using OKX.Net.Clients;
var symbol = new SharedSymbol(TradingMode.Spot, "ETH", "USDT");
var binanceSpotRestClient = new BinanceRestClient().SpotApi.SharedClient;
var okxSpotRestClient = new OKXRestClient().UnifiedApi.SharedClient;
var bitmartSpotRestClient = new BitMartRestClient().SpotApi.SharedClient;
var binanceSpotSocketClient = new BinanceSocketClient().SpotApi.SharedClient;
var okxSpotSocketClient = new OKXSocketClient().UnifiedApi.SharedClient;
var bitmartSpotSocketClient = new BitMartSocketClient().SpotApi.SharedClient;
await GetLastTradePriceAsync(binanceSpotRestClient, symbol);
await GetLastTradePriceAsync(okxSpotRestClient, symbol);
await GetLastTradePriceAsync(bitmartSpotRestClient, symbol);
Console.WriteLine();
Console.WriteLine("Press enter to start websocket");
Console.ReadLine();
await SubscribeTickerUpdatesAsync(binanceSpotSocketClient, symbol);
await SubscribeTickerUpdatesAsync(okxSpotSocketClient, symbol);
await SubscribeTickerUpdatesAsync(bitmartSpotSocketClient, symbol);
Console.ReadLine();
async Task GetLastTradePriceAsync(ISpotTickerRestClient client, SharedSymbol symbol)
{
var result = await client.GetSpotTickerAsync(new GetTickerRequest(symbol));
if (!result.Success)
{
Console.WriteLine($"Failed to get ticker: {result.Error}");
return;
}
Console.WriteLine($"{client.Exchange} {result.Data.Symbol}: {result.Data.LastPrice}");
}
async Task SubscribeTickerUpdatesAsync(ITickerSocketClient client, SharedSymbol symbol)
{
var result = await client.SubscribeToTickerUpdatesAsync(new SubscribeTickerRequest(symbol), update =>
{
Console.WriteLine($"{client.Exchange} {update.Data.Symbol} {update.Data.LastPrice}");
});
if (!result.Success)
Console.WriteLine($"Failed to subscribe ticker: {result.Error}");
}
@@ -0,0 +1,16 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net8.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Binance.Net" Version="10.7.0" />
<PackageReference Include="BitMart.Net" Version="1.4.0" />
<PackageReference Include="JK.OKX.Net" Version="2.6.0" />
</ItemGroup>
</Project>
+29 -2
View File
@@ -18,8 +18,10 @@ The following API's are directly supported. Note that there are 3rd party implem
|Bitget|[JKorf/Bitget.Net](https://github.com/JKorf/Bitget.Net)|[![Nuget version](https://img.shields.io/nuget/v/JK.Bitget.net.svg?style=flat-square)](https://www.nuget.org/packages/JK.Bitget.Net)|
|BitMart|[JKorf/BitMart.Net](https://github.com/JKorf/BitMart.Net)|[![Nuget version](https://img.shields.io/nuget/v/BitMart.net.svg?style=flat-square)](https://www.nuget.org/packages/BitMart.Net)|
|Bybit|[JKorf/Bybit.Net](https://github.com/JKorf/Bybit.Net)|[![Nuget version](https://img.shields.io/nuget/v/Bybit.net.svg?style=flat-square)](https://www.nuget.org/packages/Bybit.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)|
|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)|
|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)|
@@ -34,7 +36,7 @@ Any of these can be installed independently or install [CryptoClients.Net](https
A Discord server is available [here](https://discord.gg/MSpeEtSY8t). Feel free to join for discussion and/or questions around the CryptoExchange.Net and implementation libraries.
## Support the project
I develop and maintain this package on my own for free in my spare time, any support is greatly appreciated.
Any support is greatly appreciated.
### Donate
Make a one time donation in a crypto currency of your choice. If you prefer to donate a currency not listed here please contact me.
@@ -47,8 +49,33 @@ Make a one time donation in a crypto currency of your choice. If you prefer to d
Alternatively, sponsor me on Github using [Github Sponsors](https://github.com/sponsors/JKorf).
## Release notes
* Version 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
* Updated dependency versions, including System.Text.Json from 8.0.4 to 8.0.5 containing a vulnerability fix
* Version 8.0.1 - 07 Oct 2024
* Added cached library version properties on base client
* Added support for derserializing 0001-01-01 as datetime null value
* Added ToRfc3339String extension method for DateTime type
* Version 8.0.0 - 27 Sep 2024
* Added new cross exchange interfaces implementation
* Added new cross exchange interfaces implementation
* Supports REST, WebSocket, Spot and Futures API's
* Added various client interfaces for specific functionality
* Added SharedSymbol type, taking care of symbol formatting for different exchanges
+1046 -12
View File
File diff suppressed because it is too large Load Diff