mirror of
https://github.com/JKorf/CryptoExchange.Net.git
synced 2026-08-11 16:32:57 +00:00
e823114623
* Result types: * (Web)CallResult types are replaced by HttpResult, WebSocketResult and QueryResult with the same logic * Updated result types to record type * Result creation can be done with (Http/WebSocket/Query)Result.Ok(..) and .Fail(..) * Removed implicit result type conversion to bool, `if (result)` no longer works, instead use `if (result.Success)` * Replaced CallResult.SuccessResult with CallResult.Ok() * Fixed result object nullability hinting, for example Data might be null if Success isn't checked for true * Parameters & serialization: * Added support for `enabled` and `disabled` strings to bool converter * Removed ParameterCollection type, has been replaced by Parameters type * Removed ArraySerialization, OrderParameters and ParameterOrderComparer properties from RestApiClient, moved to ParameterSerializationsSettings * Updated RestRequestConfiguration in AuthenticationProvider.ProcessRequest to contain the full RequestDefinition instead of copied fields * Clients: * Updated Api client constructor logging parameter from ILogger to ILoggerFactory? * Added Api client constructor exchange name parameter * Added ToString overrides on base API types * Added Exchange property on BaseApiClient * Added ApiCredentials property on IRestApiClient and ISocketApiClient interfaces * Updated ILogger source from client name to topic specific client name * Removed logging from client creation * Fixed BaseRestClient SetApiCredentials not marked as virtual * Rest: * Added BaseAddress to RequestDefinition object * Updated RestApiClient AuthenticationProvider logic from private to protected and virtual * Removed RestApiClient.SendAsync baseAddress parameter removed * Removed RestApiClient.SendAsync without type parameter * WebSocket: * Updated MessageRouting definition into CreateForEvent for subscriptions and CreateForQuery for queries * Improved Query type safety with CeateForQuery which allows second parameter for specifying the result type * Renamed MessageRouter.CreateWithoutHandler to CreateVoid * Updated SocketApiClient.GetSocketConnection to check connection uri instead of Tag for finding compatible connections * Removed unused UnhandledMessageExpected property SocketApiClient * Fixed issue in SocketApiClient.GetSocketConnection causing requests to always wait the full max 10 seconds when there was a reconnecting socket * Shared APIs: * Updated Option definitions to always require the exchange name as first parameter * Added missing dedicated option types * Added Discover method on ISharedClient interface, returning info on supported capabilities and operations * Added SharedRequest GetParamValue helper method accepting multiple parameter names * Added ResetStaticExchangeParameters method on ExchangeParameters * Added Status property to SharedWithdrawal model * Added TradingModes property to SharedBalance model * Updated ExchangeSymbolCache to support multiple environments and additional key separation * Updated Shared ExchangeParameters parameter names to be case insensitive * Updated code comments * Replaced ExchangeResult with ExchangeCallResult type * Removed AsExchangeResult/ExchangeWebResult * Removed TradingMode from the response model, only maintained on models where it makes sense * Removed IListenKey support, listen keys now rely on internal management with TokenManager * Rate limiting: * Fixed websocket connection attempts counting towards rate limit even when server could not be reached * Removed host from rate limit methods, now part of the already provided RequestDefinition * Added amount parameter to RateLimit Reset method to allow partially resetting the limit * Added TokenManager implementation for automatic listenkey/token management * Added UserClientProvider base class * Added async streaming on UserDataTracker items with StreamUpdatesAsync * Added cancellation token support to UserDataTracker starting * Added Unit type for non-result types * Added ServerError constructor taking ErrorType and message to make it easier to create * Added SupportedEnvironments property to PlatformInfo * Updated SymbolOrderBook DoResyncAsync to return CallResult instead of CallResult<bool> which was redundant * Various small performance improvements
218 lines
8.8 KiB
C#
218 lines
8.8 KiB
C#
using CryptoExchange.Net.Logging.Extensions;
|
|
using CryptoExchange.Net.Objects;
|
|
using CryptoExchange.Net.RateLimiting.Guards;
|
|
using CryptoExchange.Net.RateLimiting.Interfaces;
|
|
using Microsoft.Extensions.Logging;
|
|
using System;
|
|
using System.Collections.Concurrent;
|
|
using System.Collections.Generic;
|
|
using System.Linq;
|
|
using System.Threading;
|
|
using System.Threading.Tasks;
|
|
|
|
namespace CryptoExchange.Net.RateLimiting
|
|
{
|
|
/// <inheritdoc />
|
|
public class RateLimitGate : IRateLimitGate
|
|
{
|
|
private readonly ConcurrentBag<IRateLimitGuard> _guards;
|
|
private readonly SemaphoreSlim _semaphore;
|
|
private readonly string _name;
|
|
|
|
private int _waitingCount;
|
|
|
|
/// <inheritdoc />
|
|
public event Action<RateLimitEvent>? RateLimitTriggered;
|
|
/// <inheritdoc />
|
|
public event Action<RateLimitUpdateEvent>? RateLimitUpdated;
|
|
|
|
/// <summary>
|
|
/// ctor
|
|
/// </summary>
|
|
public RateLimitGate(string name)
|
|
{
|
|
_name = name;
|
|
_guards = new ConcurrentBag<IRateLimitGuard>();
|
|
_semaphore = new SemaphoreSlim(1, 1);
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async ValueTask<CallResult> ProcessAsync(ILogger logger, int itemId, RateLimitItemType type, RequestDefinition definition, string? apiKey, int requestWeight, RateLimitingBehaviour rateLimitingBehaviour, string? keySuffix, CancellationToken ct)
|
|
{
|
|
await _semaphore.WaitAsync(ct).ConfigureAwait(false);
|
|
bool release = true;
|
|
_waitingCount++;
|
|
try
|
|
{
|
|
return await CheckGuardsAsync(_guards, logger, itemId, type, definition, apiKey, requestWeight, rateLimitingBehaviour, keySuffix, ct).ConfigureAwait(false);
|
|
}
|
|
catch (TaskCanceledException tce)
|
|
{
|
|
// The semaphore has already been released if the task was cancelled
|
|
release = false;
|
|
return CallResult.Fail(new CancellationRequestedError(tce));
|
|
}
|
|
finally
|
|
{
|
|
_waitingCount--;
|
|
if (release)
|
|
_semaphore.Release();
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async ValueTask<CallResult> ProcessSingleAsync(
|
|
ILogger logger,
|
|
int itemId,
|
|
IRateLimitGuard guard,
|
|
RateLimitItemType type,
|
|
RequestDefinition definition,
|
|
string? apiKey,
|
|
int requestWeight,
|
|
RateLimitingBehaviour rateLimitingBehaviour,
|
|
string? keySuffix,
|
|
CancellationToken ct)
|
|
{
|
|
await _semaphore.WaitAsync(ct).ConfigureAwait(false);
|
|
bool release = true;
|
|
_waitingCount++;
|
|
try
|
|
{
|
|
return await CheckGuardsAsync(new IRateLimitGuard[] { guard }, logger, itemId, type, definition, apiKey, requestWeight, rateLimitingBehaviour, keySuffix, ct).ConfigureAwait(false);
|
|
}
|
|
catch (TaskCanceledException tce)
|
|
{
|
|
// The semaphore has already been released if the task was cancelled
|
|
release = false;
|
|
return CallResult.Fail(new CancellationRequestedError(tce));
|
|
}
|
|
finally
|
|
{
|
|
_waitingCount--;
|
|
if (release)
|
|
_semaphore.Release();
|
|
}
|
|
}
|
|
|
|
private async ValueTask<CallResult> CheckGuardsAsync(IEnumerable<IRateLimitGuard> guards, ILogger logger, int itemId, RateLimitItemType type, RequestDefinition definition, string? apiKey, int requestWeight, RateLimitingBehaviour rateLimitingBehaviour, string? keySuffix, CancellationToken ct)
|
|
{
|
|
foreach (var guard in guards)
|
|
{
|
|
// Check if a wait is needed for this guard
|
|
var result = guard.Check(type, definition, apiKey, requestWeight, keySuffix);
|
|
if (result.Delay != TimeSpan.Zero && rateLimitingBehaviour == RateLimitingBehaviour.Fail)
|
|
{
|
|
// Delay is needed and limit behaviour is to fail the request
|
|
if (type == RateLimitItemType.Connection)
|
|
logger.RateLimitConnectionFailed(itemId, guard.Name, guard.Description);
|
|
else
|
|
logger.RateLimitRequestFailed(itemId, definition.Path, guard.Name, guard.Description);
|
|
|
|
RateLimitTriggered?.Invoke(new RateLimitEvent(itemId, _name, guard.Description, definition, result.Current, requestWeight, result.Limit, result.Period, result.Delay, rateLimitingBehaviour));
|
|
return CallResult.Fail(new ClientRateLimitError($"Rate limit check failed on guard {guard.Name}; {guard.Description}"));
|
|
}
|
|
|
|
if (result.Delay != TimeSpan.Zero)
|
|
{
|
|
// Delay is needed and limit behaviour is to wait for the request to be under the limit
|
|
_semaphore.Release();
|
|
|
|
var description = result.Limit == null ? guard.Description : $"{guard.Description}, Request weight: {requestWeight}, Current: {result.Current}, Limit: {result.Limit}, requests now being limited: {_waitingCount}";
|
|
if (type == RateLimitItemType.Connection)
|
|
logger.RateLimitDelayingConnection(itemId, result.Delay, guard.Name, description);
|
|
else
|
|
logger.RateLimitDelayingRequest(itemId, definition.Path, result.Delay, guard.Name, description);
|
|
|
|
RateLimitTriggered?.Invoke(new RateLimitEvent(itemId, _name, guard.Description, definition, result.Current, requestWeight, result.Limit, result.Period, result.Delay, rateLimitingBehaviour));
|
|
await Task.Delay((int)result.Delay.TotalMilliseconds + 1, ct).ConfigureAwait(false);
|
|
await _semaphore.WaitAsync(ct).ConfigureAwait(false);
|
|
return await CheckGuardsAsync(guards, logger, itemId, type, definition, apiKey, requestWeight, rateLimitingBehaviour, keySuffix, ct).ConfigureAwait(false);
|
|
}
|
|
}
|
|
|
|
// Apply the weight on each guard
|
|
foreach (var guard in guards)
|
|
{
|
|
var result = guard.ApplyWeight(type, definition, apiKey, requestWeight, keySuffix);
|
|
if (result.IsApplied)
|
|
{
|
|
RateLimitUpdated?.Invoke(new RateLimitUpdateEvent(itemId, _name, guard.Description, result.Current, result.Limit, result.Period));
|
|
|
|
if (logger.IsEnabled(LogLevel.Trace))
|
|
{
|
|
if (type == RateLimitItemType.Connection)
|
|
logger.RateLimitAppliedConnection(itemId, guard.Name, guard.Description, result.Current);
|
|
else
|
|
logger.RateLimitAppliedRequest(itemId, definition.Path, guard.Name, guard.Description, result.Current);
|
|
}
|
|
}
|
|
}
|
|
|
|
return CallResult.Ok();
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public IRateLimitGate AddGuard(IRateLimitGuard guard)
|
|
{
|
|
_guards.Add(guard);
|
|
return this;
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task SetRetryAfterGuardAsync(DateTime retryAfter, RateLimitItemType type)
|
|
{
|
|
await _semaphore.WaitAsync().ConfigureAwait(false);
|
|
|
|
try
|
|
{
|
|
var retryAfterGuard = _guards.OfType<RetryAfterGuard>().SingleOrDefault();
|
|
if (retryAfterGuard == null)
|
|
_guards.Add(new RetryAfterGuard(retryAfter, type));
|
|
else
|
|
retryAfterGuard.UpdateAfter(retryAfter);
|
|
}
|
|
finally
|
|
{
|
|
_semaphore.Release();
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task<DateTime?> GetRetryAfterTime()
|
|
{
|
|
await _semaphore.WaitAsync().ConfigureAwait(false);
|
|
try
|
|
{
|
|
var retryAfterGuard = _guards.OfType<RetryAfterGuard>().SingleOrDefault();
|
|
return retryAfterGuard?.After;
|
|
}
|
|
finally
|
|
{
|
|
_semaphore.Release();
|
|
}
|
|
}
|
|
|
|
|
|
/// <inheritdoc />
|
|
public async Task ResetAsync(
|
|
RateLimitItemType type,
|
|
RequestDefinition definition,
|
|
string? apiKey,
|
|
string? keySuffix,
|
|
int? amount,
|
|
CancellationToken ct)
|
|
{
|
|
await _semaphore.WaitAsync(ct).ConfigureAwait(false);
|
|
try
|
|
{
|
|
foreach (var guard in _guards)
|
|
guard.Reset(type, definition, apiKey, keySuffix, amount);
|
|
}
|
|
finally
|
|
{
|
|
_semaphore.Release();
|
|
}
|
|
}
|
|
}
|
|
}
|