mirror of
https://github.com/JKorf/CryptoExchange.Net.git
synced 2026-08-12 00:43:03 +00:00
Compare commits
23 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a0e588c3de | |||
| 9e86a08327 | |||
| ed007b5272 | |||
| bdd5526244 | |||
| ce35e30688 | |||
| b40f72b1b0 | |||
| 31a6cf285b | |||
| 1842f4fda0 | |||
| 7a58902ab6 | |||
| 3cb91296ca | |||
| 130ed40580 | |||
| 94cb2caf0b | |||
| 917d060827 | |||
| c58bc2be07 | |||
| ff3356e2b4 | |||
| 79434c7be5 | |||
| 168dabc11f | |||
| 71ee263683 | |||
| 7239b9c289 | |||
| 84d36544e4 | |||
| a71f57ae7f | |||
| 6e5bcd5e9a | |||
| 4131e563c3 |
@@ -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; }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
@@ -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++;
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
{
|
||||
|
||||
@@ -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.1</PackageVersion>
|
||||
<AssemblyVersion>8.0.1</AssemblyVersion>
|
||||
<FileVersion>8.0.1</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>
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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())
|
||||
{
|
||||
|
||||
@@ -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>
|
||||
|
||||
@@ -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),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5,19 +5,21 @@
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Binance.Net" Version="10.5.0" />
|
||||
<PackageReference Include="Bitfinex.Net" Version="7.8.0" />
|
||||
<PackageReference Include="BitMart.Net" Version="1.2.0" />
|
||||
<PackageReference Include="Bybit.Net" Version="3.14.0" />
|
||||
<PackageReference Include="CoinEx.Net" Version="7.7.0" />
|
||||
<PackageReference Include="GateIo.Net" Version="1.6.0" />
|
||||
<PackageReference Include="JK.BingX.Net" Version="1.11.0" />
|
||||
<PackageReference Include="JK.Bitget.Net" Version="1.10.0" />
|
||||
<PackageReference Include="JK.Mexc.Net" Version="1.8.0" />
|
||||
<PackageReference Include="JK.OKX.Net" Version="2.4.0" />
|
||||
<PackageReference Include="JKorf.HTX.Net" Version="6.1.0" />
|
||||
<PackageReference Include="KrakenExchange.Net" Version="4.12.0" />
|
||||
<PackageReference Include="Kucoin.Net" Version="5.14.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>
|
||||
|
||||
|
||||
@@ -5,7 +5,9 @@
|
||||
@inject IBitMartRestClient bitmartClient
|
||||
@inject IBitgetRestClient bitgetClient
|
||||
@inject IBybitRestClient bybitClient
|
||||
@inject ICoinbaseRestClient coinbaseClient
|
||||
@inject ICoinExRestClient coinexClient
|
||||
@inject ICryptoComRestClient cryptocomClient
|
||||
@inject IGateIoRestClient gateioClient
|
||||
@inject IHTXRestClient huobiClient
|
||||
@inject IKrakenRestClient krakenClient
|
||||
@@ -30,7 +32,9 @@
|
||||
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 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");
|
||||
@@ -58,9 +62,15 @@
|
||||
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 (cryptocomTask.Result.Success)
|
||||
_prices.Add("CryptoCom", cryptocomTask.Result.Data.First().LastPrice ?? 0);
|
||||
|
||||
if (gateioTask.Result.Success)
|
||||
_prices.Add("GateIo", gateioTask.Result.Data.First().LastPrice);
|
||||
|
||||
|
||||
@@ -5,7 +5,9 @@
|
||||
@inject IBitgetSocketClient bitgetSocketClient
|
||||
@inject IBitMartSocketClient bitmartSocketClient
|
||||
@inject IBybitSocketClient bybitSocketClient
|
||||
@inject ICoinbaseSocketClient coinbaseSocketClient
|
||||
@inject ICoinExSocketClient coinExSocketClient
|
||||
@inject ICryptoComSocketClient cryptocomSocketClient
|
||||
@inject IGateIoSocketClient gateioSocketClient
|
||||
@inject IHTXSocketClient htxSocketClient
|
||||
@inject IKrakenSocketClient krakenSocketClient
|
||||
@@ -39,9 +41,11 @@
|
||||
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)),
|
||||
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.LastTrade.Price)),
|
||||
krakenSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH/XBT", data => UpdateData("Kraken", data.Data.LastPrice)),
|
||||
kucoinSocketClient.SpotApi.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("Kucoin", data.Data.LastPrice ?? 0)),
|
||||
mexcSocketClient.SpotApi.SubscribeToMiniTickerUpdatesAsync("ETHBTC", data => UpdateData("Mexc", data.Data.LastPrice)),
|
||||
okxSocketClient.UnifiedApi.ExchangeData.SubscribeToTickerUpdatesAsync("ETH-BTC", data => UpdateData("OKX", data.Data.LastPrice ?? 0)),
|
||||
|
||||
@@ -8,7 +8,9 @@
|
||||
@using BitMart.Net.Interfaces;
|
||||
@using Bybit.Net.Interfaces
|
||||
@using CoinEx.Net.Interfaces
|
||||
@using Coinbase.Net.Interfaces
|
||||
@using CryptoExchange.Net.Interfaces
|
||||
@using CryptoCom.Net.Interfaces
|
||||
@using GateIo.Net.Interfaces
|
||||
@using HTX.Net.Interfaces
|
||||
@using Kraken.Net.Interfaces
|
||||
@@ -22,7 +24,9 @@
|
||||
@inject IBitgetOrderBookFactory bitgetFactory
|
||||
@inject IBitMartOrderBookFactory bitmartFactory
|
||||
@inject IBybitOrderBookFactory bybitFactory
|
||||
@inject ICoinbaseOrderBookFactory coinbaseFactory
|
||||
@inject ICoinExOrderBookFactory coinExFactory
|
||||
@inject ICryptoComOrderBookFactory cryptocomFactory
|
||||
@inject IGateIoOrderBookFactory gateioFactory
|
||||
@inject IHTXOrderBookFactory htxFactory
|
||||
@inject IKrakenOrderBookFactory krakenFactory
|
||||
@@ -68,7 +72,9 @@
|
||||
{ "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") },
|
||||
{ "CryptoCom", cryptocomFactory.CreateExchange("ETH_BTC") },
|
||||
{ "GateIo", gateioFactory.CreateSpot("ETH_BTC") },
|
||||
{ "HTX", htxFactory.CreateSpot("ethbtc") },
|
||||
{ "Kraken", krakenFactory.CreateSpot("ETH/XBT") },
|
||||
|
||||
@@ -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">
|
||||
|
||||
@@ -41,7 +41,9 @@ namespace BlazorClient
|
||||
services.AddBitget();
|
||||
services.AddBitMart();
|
||||
services.AddBybit();
|
||||
services.AddCoinbase();
|
||||
services.AddCoinEx();
|
||||
services.AddCryptoCom();
|
||||
services.AddGateIo();
|
||||
services.AddHTX();
|
||||
services.AddKraken();
|
||||
|
||||
@@ -14,7 +14,9 @@
|
||||
@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 CryptoCom.Net.Interfaces.Clients;
|
||||
@using GateIo.Net.Interfaces.Clients;
|
||||
@using HTX.Net.Interfaces.Clients;
|
||||
@using Kraken.Net.Interfaces.Clients;
|
||||
|
||||
@@ -6,18 +6,20 @@
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Binance.Net" Version="10.5.0" />
|
||||
<PackageReference Include="Bitfinex.Net" Version="7.8.0" />
|
||||
<PackageReference Include="BitMart.Net" Version="1.2.0" />
|
||||
<PackageReference Include="Bybit.Net" Version="3.14.0" />
|
||||
<PackageReference Include="CoinEx.Net" Version="7.7.0" />
|
||||
<PackageReference Include="GateIo.Net" Version="1.6.0" />
|
||||
<PackageReference Include="JK.Bitget.Net" Version="1.10.0" />
|
||||
<PackageReference Include="JK.Mexc.Net" Version="1.8.0" />
|
||||
<PackageReference Include="JK.OKX.Net" Version="2.4.0" />
|
||||
<PackageReference Include="JKorf.HTX.Net" Version="6.1.0" />
|
||||
<PackageReference Include="KrakenExchange.Net" Version="4.12.0" />
|
||||
<PackageReference Include="Kucoin.Net" Version="5.14.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()
|
||||
|
||||
@@ -8,9 +8,9 @@
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Binance.Net" Version="10.5.0" />
|
||||
<PackageReference Include="BitMart.Net" Version="1.2.0" />
|
||||
<PackageReference Include="JK.OKX.Net" Version="2.4.0" />
|
||||
<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>
|
||||
|
||||
@@ -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)|[](https://www.nuget.org/packages/JK.Bitget.Net)|
|
||||
|BitMart|[JKorf/BitMart.Net](https://github.com/JKorf/BitMart.Net)|[](https://www.nuget.org/packages/BitMart.Net)|
|
||||
|Bybit|[JKorf/Bybit.Net](https://github.com/JKorf/Bybit.Net)|[](https://www.nuget.org/packages/Bybit.Net)|
|
||||
|Coinbase|[JKorf/Coinbase.Net](https://github.com/JKorf/Coinbase.Net)|[](https://www.nuget.org/packages/JKorf.Coinbase.Net)|
|
||||
|CoinEx|[JKorf/CoinEx.Net](https://github.com/JKorf/CoinEx.Net)|[](https://www.nuget.org/packages/CoinEx.Net)|
|
||||
|CoinGecko|[JKorf/CoinGecko.Net](https://github.com/JKorf/CoinGecko.Net)|[](https://www.nuget.org/packages/CoinGecko.Net)|
|
||||
|Crypto.com|[JKorf/CryptoCom.Net](https://github.com/JKorf/CryptoCom.Net)|[](https://www.nuget.org/packages/CryptoCom.Net)|
|
||||
|Gate.io|[JKorf/GateIo.Net](https://github.com/JKorf/GateIo.Net)|[](https://www.nuget.org/packages/GateIo.Net)|
|
||||
|HTX|[JKorf/HTX.Net](https://github.com/JKorf/HTX.Net)|[](https://www.nuget.org/packages/JKorf.HTX.Net)|
|
||||
|Kraken|[JKorf/Kraken.Net](https://github.com/JKorf/Kraken.Net)|[](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,13 +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
|
||||
|
||||
+1025
-10
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user