1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-10-04 02:11:11 +00:00

Compare commits

...

3 Commits

10 changed files with 153 additions and 10 deletions
@@ -533,9 +533,6 @@ namespace CryptoExchange.Net.Clients
if (!connectResult.Success) if (!connectResult.Success)
return connectResult; return connectResult;
if (ClientOptions.DelayAfterConnect != TimeSpan.Zero)
await Task.Delay(ClientOptions.DelayAfterConnect).ConfigureAwait(false);
if (!authenticated || socket.Authenticated) if (!authenticated || socket.Authenticated)
return CallResult.Ok(); return CallResult.Ok();
@@ -862,7 +859,7 @@ namespace CryptoExchange.Net.Clients
RateLimitAdmissionCallbackRequest = () => AdmissionOverride.Value, RateLimitAdmissionCallbackRequest = () => AdmissionOverride.Value,
Proxy = ClientOptions.Proxy, Proxy = ClientOptions.Proxy,
Timeout = ApiOptions.SocketNoDataTimeout ?? ClientOptions.SocketNoDataTimeout, Timeout = ApiOptions.SocketNoDataTimeout ?? ClientOptions.SocketNoDataTimeout,
ReceiveBufferSize = ClientOptions.ReceiveBufferSize, ReceiveBufferSize = ClientOptions.ReceiveBufferSize
}; };
/// <summary> /// <summary>
@@ -19,6 +19,11 @@ namespace CryptoExchange.Net.Objects.Options
/// </summary> /// </summary>
public int? MaxSocketConnections { get; set; } public int? MaxSocketConnections { get; set; }
/// <summary>
/// The time to wait after connecting a socket before sending messages. Can be used for API's which will rate limit if you subscribe directly after connecting.
/// </summary>
public TimeSpan? DelayAfterConnect { get; set; }
/// <summary> /// <summary>
/// Set the values of this options on the target options /// Set the values of this options on the target options
/// </summary> /// </summary>
@@ -28,6 +33,7 @@ namespace CryptoExchange.Net.Objects.Options
item.SocketNoDataTimeout = SocketNoDataTimeout; item.SocketNoDataTimeout = SocketNoDataTimeout;
item.AutoTimestamp = AutoTimestamp; item.AutoTimestamp = AutoTimestamp;
item.MaxSocketConnections = MaxSocketConnections; item.MaxSocketConnections = MaxSocketConnections;
item.DelayAfterConnect = DelayAfterConnect;
return item; return item;
} }
} }
@@ -0,0 +1,25 @@
using CryptoExchange.Net.Interfaces;
namespace CryptoExchange.Net.SharedApis
{
/// <summary>
/// Order book info
/// </summary>
public record SharedIncrementalOrderBook : SharedOrderBook
{
/// <summary>
/// The sequence number of the first book update in this update.
/// </summary>
public long? StartSequenceNumber { get; set; }
/// <summary>
/// ctor
/// </summary>
public SharedIncrementalOrderBook(SharedQuantityType quantityType, long? startSequenceNumber, long? endSequenceNumber, ISymbolOrderBookEntry[] asks, ISymbolOrderBookEntry[] bids)
:base(quantityType, endSequenceNumber, asks, bids)
{
StartSequenceNumber = startSequenceNumber;
}
}
}
+16 -2
View File
@@ -668,6 +668,20 @@ namespace CryptoExchange.Net.SharedApis
.ParallelEnumerateAsync(); .ParallelEnumerateAsync();
} }
/// <summary>
/// Subscribe to funding info updates for all capabilities in parallel and return results as they arrive
/// </summary>
public static IAsyncEnumerable<WebSocketResult<UpdateSubscription>> SubscribeAllAsync(
this IEnumerable<SharedCapabilityResolution<ISubscribeFundingInfoSocket>> capabilities,
SubscribeFundingInfoRequest request,
Action<DataEvent<SharedFundingInfo>> onData,
CancellationToken ct = default)
{
return capabilities
.Select(x => x.Capability.SubscribeToFundingInfoUpdatesAsync(request, onData, ct))
.ParallelEnumerateAsync();
}
/// <summary> /// <summary>
/// Subscribe to kline updates for all capabilities in parallel and return results as they arrive /// Subscribe to kline updates for all capabilities in parallel and return results as they arrive
/// </summary> /// </summary>
@@ -716,11 +730,11 @@ namespace CryptoExchange.Net.SharedApis
public static IAsyncEnumerable<WebSocketResult<UpdateSubscription>> SubscribeAllAsync( public static IAsyncEnumerable<WebSocketResult<UpdateSubscription>> SubscribeAllAsync(
this IEnumerable<SharedCapabilityResolution<ISubscribeIncrementalOrderBookSocket>> capabilities, this IEnumerable<SharedCapabilityResolution<ISubscribeIncrementalOrderBookSocket>> capabilities,
SubscribeOrderBookRequest request, SubscribeOrderBookRequest request,
Action<DataEvent<SharedOrderBook>> onData, Action<DataEvent<SharedIncrementalOrderBook>> onData,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return capabilities return capabilities
.Select(x => x.Capability.SubscribeToOrderBookUpdatesAsync(request, onData, ct)) .Select(x => x.Capability.SubscribeToIncrementalOrderBookUpdatesAsync(request, onData, ct))
.ParallelEnumerateAsync(); .ParallelEnumerateAsync();
} }
@@ -0,0 +1,28 @@
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.Objects.Sockets;
using System;
using System.Threading;
using System.Threading.Tasks;
namespace CryptoExchange.Net.SharedApis
{
/// <summary>
/// Operation for subscribing to funding info updates
/// </summary>
public interface ISubscribeFundingInfoSocket : ISharedSubscription
{
/// <summary>
/// Funding info subscription options
/// </summary>
SubscribeFundingInfoOptions SubscribeFundingInfoOptions { get; }
/// <summary>
/// Subscribe to funding info updates
/// </summary>
/// <param name="request">Request info</param>
/// <param name="handler">Update handler</param>
/// <param name="ct">Cancellation token, can be used to stop the updates</param>
/// <returns></returns>
Task<WebSocketResult<UpdateSubscription>> SubscribeToFundingInfoUpdatesAsync(SubscribeFundingInfoRequest request, Action<DataEvent<SharedFundingInfo>> handler, CancellationToken ct = default);
}
}
@@ -0,0 +1,29 @@
using CryptoExchange.Net.Objects;
using System;
using System.Linq;
namespace CryptoExchange.Net.SharedApis
{
/// <summary>
/// Options for subscribing to funding info updates
/// </summary>
public class SubscribeFundingInfoOptions : CapabilityOptions<SubscribeFundingInfoRequest, ISubscribeFundingInfoSocket>
{
/// <inheritdoc />
public override string Description => "Subscribe to funding info updates for a symbol";
private static readonly RequestParameterDescription[] _defaultParameterRules = new[]
{
RequestParameterRule<SubscribeFundingInfoRequest>.Optional(x => x.Symbol, "The symbol to subscribe to", new SharedSymbol(TradingMode.PerpetualLinear, "ETH", "USDT")),
RequestParameterRule<SubscribeFundingInfoRequest>.Optional(x => x.Symbols, "The symbols to subscribe to", new[] { new SharedSymbol(TradingMode.PerpetualLinear, "ETH", "USDT") }),
};
/// <summary>
/// ctor
/// </summary>
public SubscribeFundingInfoOptions(string exchange, bool needsAuthentication)
: base(exchange, needsAuthentication, nameof(ISubscribeFundingInfoSocket.SubscribeToFundingInfoUpdatesAsync), _defaultParameterRules, SharedTradingModeSets.Futures)
{
}
}
}
@@ -0,0 +1,31 @@
using System;
using System.Collections.Generic;
namespace CryptoExchange.Net.SharedApis
{
/// <summary>
/// Request to subscribe to funding info updates for a symbol
/// </summary>
public record SubscribeFundingInfoRequest : SharedSymbolRequest
{
/// <summary>
/// ctor
/// </summary>
/// <param name="symbol">The symbol to subscribe to</param>
/// <param name="exchangeParameters">Exchange specific parameters</param>
public SubscribeFundingInfoRequest(SharedSymbol symbol, ExchangeParameters? exchangeParameters = null)
: base(symbol, exchangeParameters)
{
}
/// <summary>
/// ctor
/// </summary>
/// <param name="symbols">The symbols to subscribe to</param>
/// <param name="exchangeParameters">Exchange specific parameters</param>
public SubscribeFundingInfoRequest(IEnumerable<SharedSymbol> symbols, ExchangeParameters? exchangeParameters = null)
: base(symbols, exchangeParameters)
{
}
}
}
@@ -14,7 +14,7 @@ namespace CryptoExchange.Net.SharedApis
/// <summary> /// <summary>
/// Order book subscription options /// Order book subscription options
/// </summary> /// </summary>
SubscribeOrderBookOptions SubscribeOrderBookOptions { get; } SubscribeIncrementalOrderBookOptions SubscribeIncrementalOrderBookOptions { get; }
/// <summary> /// <summary>
/// Subscribe to incremental order book updates for a symbol /// Subscribe to incremental order book updates for a symbol
@@ -23,6 +23,6 @@ namespace CryptoExchange.Net.SharedApis
/// <param name="handler">Update handler</param> /// <param name="handler">Update handler</param>
/// <param name="ct">Cancellation token, can be used to stop the updates</param> /// <param name="ct">Cancellation token, can be used to stop the updates</param>
/// <returns></returns> /// <returns></returns>
Task<WebSocketResult<UpdateSubscription>> SubscribeToOrderBookUpdatesAsync(SubscribeOrderBookRequest request, Action<DataEvent<SharedOrderBook>> handler, CancellationToken ct = default); Task<WebSocketResult<UpdateSubscription>> SubscribeToIncrementalOrderBookUpdatesAsync(SubscribeOrderBookRequest request, Action<DataEvent<SharedIncrementalOrderBook>> handler, CancellationToken ct = default);
} }
} }
@@ -33,7 +33,7 @@ namespace CryptoExchange.Net.SharedApis
/// ctor /// ctor
/// </summary> /// </summary>
public SubscribeIncrementalOrderBookOptions(string exchange, bool needsAuthentication, int[] limits, SharedOrderBookSubscriptionType updateType) public SubscribeIncrementalOrderBookOptions(string exchange, bool needsAuthentication, int[] limits, SharedOrderBookSubscriptionType updateType)
: base(exchange, needsAuthentication, nameof(ISubscribeIncrementalOrderBookSocket.SubscribeToOrderBookUpdatesAsync), _defaultParameterRules) : base(exchange, needsAuthentication, nameof(ISubscribeIncrementalOrderBookSocket.SubscribeToIncrementalOrderBookUpdatesAsync), _defaultParameterRules)
{ {
SupportedLimits = limits; SupportedLimits = limits;
UpdateType = updateType; UpdateType = updateType;
@@ -424,6 +424,9 @@ namespace CryptoExchange.Net.Sockets.Default
{ {
try try
{ {
if ((ApiClient.ApiOptions.DelayAfterConnect ?? ApiClient.ClientOptions.DelayAfterConnect) != TimeSpan.Zero)
await Task.Delay(ApiClient.ApiOptions.DelayAfterConnect ?? ApiClient.ClientOptions.DelayAfterConnect).ConfigureAwait(false);
var reconnectSuccessful = await ProcessReconnectAsync().ConfigureAwait(false); var reconnectSuccessful = await ProcessReconnectAsync().ConfigureAwait(false);
if (!reconnectSuccessful.Success) if (!reconnectSuccessful.Success)
{ {
@@ -616,7 +619,17 @@ namespace CryptoExchange.Net.Sockets.Default
/// Connect the websocket /// Connect the websocket
/// </summary> /// </summary>
/// <returns></returns> /// <returns></returns>
public async Task<CallResult> ConnectAsync(CancellationToken ct) => await _socket.ConnectAsync(ct).ConfigureAwait(false); public async Task<CallResult> ConnectAsync(CancellationToken ct)
{
var result = await _socket.ConnectAsync(ct).ConfigureAwait(false);
if (!result.Success)
return result;
if ((ApiClient.ApiOptions.DelayAfterConnect ?? ApiClient.ClientOptions.DelayAfterConnect) != TimeSpan.Zero)
await Task.Delay(ApiClient.ApiOptions.DelayAfterConnect ?? ApiClient.ClientOptions.DelayAfterConnect).ConfigureAwait(false);
return result;
}
/// <summary> /// <summary>
/// Retrieve the underlying socket /// Retrieve the underlying socket