mirror of
https://github.com/JKorf/CryptoExchange.Net.git
synced 2026-08-14 09:52:53 +00:00
Compare commits
36 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ea9375d582 | |||
| 2cf3c93e5e | |||
| ca888d8e41 | |||
| 2040b1c175 | |||
| d451c18821 | |||
| c13dfa4461 | |||
| c2080ef75f | |||
| 6b252e8024 | |||
| d06bd5f176 | |||
| d55fc8da65 | |||
| 01184f2c5d | |||
| cadc93c2f0 | |||
| 2600a51461 | |||
| 9e6a86ba8b | |||
| c4430d63fa | |||
| f3e1cfef33 | |||
| cc3053719c | |||
| cd6907e601 | |||
| 8fe00693bd | |||
| fb90d1e015 | |||
| 4b44861e43 | |||
| e42ca4ab5a | |||
| 5b97f6dd67 | |||
| a9813ecb0a | |||
| c7069a4049 | |||
| 5683ae0b3c | |||
| 1c8cf5ac98 | |||
| ad7231ec56 | |||
| 7e4a607391 | |||
| 2d470d18e2 | |||
| cb9a766c3b | |||
| 94b8184f7b | |||
| 270ea06f24 | |||
| 536afa92da | |||
| 11c48b3341 | |||
| f514e172d7 |
@@ -6,10 +6,10 @@
|
|||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
<packagereference Include="Microsoft.NET.Test.Sdk" Version="17.1.0-preview-20211130-02"></packagereference>
|
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.1.0-preview-20211130-02"></PackageReference>
|
||||||
<PackageReference Include="Moq" Version="4.16.1" />
|
<PackageReference Include="Moq" Version="4.16.1" />
|
||||||
<packagereference Include="NUnit" Version="3.13.2"></packagereference>
|
<PackageReference Include="NUnit" Version="3.13.2"></PackageReference>
|
||||||
<packagereference Include="NUnit3TestAdapter" Version="4.2.0"></packagereference>
|
<PackageReference Include="NUnit3TestAdapter" Version="4.2.0"></PackageReference>
|
||||||
</ItemGroup>
|
</ItemGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
|
|||||||
@@ -135,7 +135,7 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
|
|||||||
throw new NotImplementedException();
|
throw new NotImplementedException();
|
||||||
}
|
}
|
||||||
|
|
||||||
protected override TimeSyncInfo GetTimeSyncInfo()
|
public override TimeSyncInfo GetTimeSyncInfo()
|
||||||
{
|
{
|
||||||
throw new NotImplementedException();
|
throw new NotImplementedException();
|
||||||
}
|
}
|
||||||
@@ -161,7 +161,7 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
|
|||||||
throw new NotImplementedException();
|
throw new NotImplementedException();
|
||||||
}
|
}
|
||||||
|
|
||||||
protected override TimeSyncInfo GetTimeSyncInfo()
|
public override TimeSyncInfo GetTimeSyncInfo()
|
||||||
{
|
{
|
||||||
throw new NotImplementedException();
|
throw new NotImplementedException();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -38,6 +38,10 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
|
|||||||
|
|
||||||
public double IncomingKbps => throw new NotImplementedException();
|
public double IncomingKbps => throw new NotImplementedException();
|
||||||
|
|
||||||
|
public Uri Uri => new Uri("");
|
||||||
|
|
||||||
|
public TimeSpan KeepAliveInterval { get; set; }
|
||||||
|
|
||||||
public static int lastId = 0;
|
public static int lastId = 0;
|
||||||
public static object lastIdLock = new object();
|
public static object lastIdLock = new object();
|
||||||
|
|
||||||
@@ -111,5 +115,11 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
|
|||||||
{
|
{
|
||||||
OnError?.Invoke(error);
|
OnError?.Invoke(error);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public async Task ProcessAsync()
|
||||||
|
{
|
||||||
|
while (Connected)
|
||||||
|
await Task.Delay(50);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -63,6 +63,7 @@ namespace CryptoExchange.Net
|
|||||||
log = new Log(name);
|
log = new Log(name);
|
||||||
log.UpdateWriters(options.LogWriters);
|
log.UpdateWriters(options.LogWriters);
|
||||||
log.Level = options.LogLevel;
|
log.Level = options.LogLevel;
|
||||||
|
options.OnLoggingChanged += HandleLogConfigChange;
|
||||||
|
|
||||||
ClientOptions = options;
|
ClientOptions = options;
|
||||||
|
|
||||||
@@ -282,12 +283,22 @@ namespace CryptoExchange.Net
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Handle a change in the client options log config
|
||||||
|
/// </summary>
|
||||||
|
private void HandleLogConfigChange()
|
||||||
|
{
|
||||||
|
log.UpdateWriters(ClientOptions.LogWriters);
|
||||||
|
log.Level = ClientOptions.LogLevel;
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Dispose
|
/// Dispose
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public virtual void Dispose()
|
public virtual void Dispose()
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Debug, "Disposing client");
|
log.Write(LogLevel.Debug, "Disposing client");
|
||||||
|
ClientOptions.OnLoggingChanged -= HandleLogConfigChange;
|
||||||
foreach (var client in ApiClients)
|
foreach (var client in ApiClients)
|
||||||
client.Dispose();
|
client.Dispose();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -61,6 +61,44 @@ namespace CryptoExchange.Net
|
|||||||
apiClient.SetApiCredentials(credentials);
|
apiClient.SetApiCredentials(credentials);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Execute a request to the uri and returns if it was successful
|
||||||
|
/// </summary>
|
||||||
|
/// <param name="apiClient">The API client the request is for</param>
|
||||||
|
/// <param name="uri">The uri to send the request to</param>
|
||||||
|
/// <param name="method">The method of the request</param>
|
||||||
|
/// <param name="cancellationToken">Cancellation token</param>
|
||||||
|
/// <param name="parameters">The parameters of the request</param>
|
||||||
|
/// <param name="signed">Whether or not the request should be authenticated</param>
|
||||||
|
/// <param name="parameterPosition">Where the parameters should be placed, overwrites the value set in the client</param>
|
||||||
|
/// <param name="arraySerialization">How array parameters should be serialized, overwrites the value set in the client</param>
|
||||||
|
/// <param name="requestWeight">Credits used for the request</param>
|
||||||
|
/// <param name="deserializer">The JsonSerializer to use for deserialization</param>
|
||||||
|
/// <param name="additionalHeaders">Additional headers to send with the request</param>
|
||||||
|
/// <param name="ignoreRatelimit">Ignore rate limits for this request</param>
|
||||||
|
/// <returns></returns>
|
||||||
|
[return: NotNull]
|
||||||
|
protected virtual async Task<WebCallResult> SendRequestAsync(RestApiClient apiClient,
|
||||||
|
Uri uri,
|
||||||
|
HttpMethod method,
|
||||||
|
CancellationToken cancellationToken,
|
||||||
|
Dictionary<string, object>? parameters = null,
|
||||||
|
bool signed = false,
|
||||||
|
HttpMethodParameterPosition? parameterPosition = null,
|
||||||
|
ArrayParametersSerialization? arraySerialization = null,
|
||||||
|
int requestWeight = 1,
|
||||||
|
JsonSerializer? deserializer = null,
|
||||||
|
Dictionary<string, string>? additionalHeaders = null,
|
||||||
|
bool ignoreRatelimit = false)
|
||||||
|
{
|
||||||
|
var request = await PrepareRequestAsync(apiClient, uri, method, cancellationToken, parameters, signed, parameterPosition, arraySerialization, requestWeight, deserializer, additionalHeaders, ignoreRatelimit).ConfigureAwait(false);
|
||||||
|
if (!request)
|
||||||
|
return new WebCallResult(request.Error!);
|
||||||
|
|
||||||
|
var result = await GetResponseAsync<object>(apiClient, request.Data, deserializer, cancellationToken, true).ConfigureAwait(false);
|
||||||
|
return result.AsDataless();
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Execute a request to the uri and deserialize the response into the provided type parameter
|
/// Execute a request to the uri and deserialize the response into the provided type parameter
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -93,16 +131,58 @@ namespace CryptoExchange.Net
|
|||||||
Dictionary<string, string>? additionalHeaders = null,
|
Dictionary<string, string>? additionalHeaders = null,
|
||||||
bool ignoreRatelimit = false
|
bool ignoreRatelimit = false
|
||||||
) where T : class
|
) where T : class
|
||||||
|
{
|
||||||
|
var request = await PrepareRequestAsync(apiClient, uri, method, cancellationToken, parameters, signed, parameterPosition, arraySerialization, requestWeight, deserializer, additionalHeaders, ignoreRatelimit).ConfigureAwait(false);
|
||||||
|
if (!request)
|
||||||
|
return new WebCallResult<T>(request.Error!);
|
||||||
|
|
||||||
|
return await GetResponseAsync<T>(apiClient, request.Data, deserializer, cancellationToken, false).ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Prepares a request to be sent to the server
|
||||||
|
/// </summary>
|
||||||
|
/// <param name="apiClient">The API client the request is for</param>
|
||||||
|
/// <param name="uri">The uri to send the request to</param>
|
||||||
|
/// <param name="method">The method of the request</param>
|
||||||
|
/// <param name="cancellationToken">Cancellation token</param>
|
||||||
|
/// <param name="parameters">The parameters of the request</param>
|
||||||
|
/// <param name="signed">Whether or not the request should be authenticated</param>
|
||||||
|
/// <param name="parameterPosition">Where the parameters should be placed, overwrites the value set in the client</param>
|
||||||
|
/// <param name="arraySerialization">How array parameters should be serialized, overwrites the value set in the client</param>
|
||||||
|
/// <param name="requestWeight">Credits used for the request</param>
|
||||||
|
/// <param name="deserializer">The JsonSerializer to use for deserialization</param>
|
||||||
|
/// <param name="additionalHeaders">Additional headers to send with the request</param>
|
||||||
|
/// <param name="ignoreRatelimit">Ignore rate limits for this request</param>
|
||||||
|
/// <returns></returns>
|
||||||
|
protected virtual async Task<CallResult<IRequest>> PrepareRequestAsync(RestApiClient apiClient,
|
||||||
|
Uri uri,
|
||||||
|
HttpMethod method,
|
||||||
|
CancellationToken cancellationToken,
|
||||||
|
Dictionary<string, object>? parameters = null,
|
||||||
|
bool signed = false,
|
||||||
|
HttpMethodParameterPosition? parameterPosition = null,
|
||||||
|
ArrayParametersSerialization? arraySerialization = null,
|
||||||
|
int requestWeight = 1,
|
||||||
|
JsonSerializer? deserializer = null,
|
||||||
|
Dictionary<string, string>? additionalHeaders = null,
|
||||||
|
bool ignoreRatelimit = false)
|
||||||
{
|
{
|
||||||
var requestId = NextId();
|
var requestId = NextId();
|
||||||
|
|
||||||
if (signed)
|
if (signed)
|
||||||
{
|
{
|
||||||
var syncTimeResult = await apiClient.SyncTimeAsync().ConfigureAwait(false);
|
var syncTask = apiClient.SyncTimeAsync();
|
||||||
if (!syncTimeResult)
|
var timeSyncInfo = apiClient.GetTimeSyncInfo();
|
||||||
|
if (timeSyncInfo.TimeSyncState.LastSyncTime == default)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Debug, $"[{requestId}] Failed to sync time, aborting request: " + syncTimeResult.Error);
|
// Initially with first request we'll need to wait for the time syncing, if it's not the first request we can just continue
|
||||||
return syncTimeResult.As<T>(default);
|
var syncTimeResult = await syncTask.ConfigureAwait(false);
|
||||||
|
if (!syncTimeResult)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Debug, $"[{requestId}] Failed to sync time, aborting request: " + syncTimeResult.Error);
|
||||||
|
return syncTimeResult.As<IRequest>(default);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -112,20 +192,20 @@ namespace CryptoExchange.Net
|
|||||||
{
|
{
|
||||||
var limitResult = await limiter.LimitRequestAsync(log, uri.AbsolutePath, method, signed, apiClient.Options.ApiCredentials?.Key, apiClient.Options.RateLimitingBehaviour, requestWeight, cancellationToken).ConfigureAwait(false);
|
var limitResult = await limiter.LimitRequestAsync(log, uri.AbsolutePath, method, signed, apiClient.Options.ApiCredentials?.Key, apiClient.Options.RateLimitingBehaviour, requestWeight, cancellationToken).ConfigureAwait(false);
|
||||||
if (!limitResult.Success)
|
if (!limitResult.Success)
|
||||||
return new WebCallResult<T>(limitResult.Error!);
|
return new CallResult<IRequest>(limitResult.Error!);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (signed && apiClient.AuthenticationProvider == null)
|
if (signed && apiClient.AuthenticationProvider == null)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"[{requestId}] Request {uri.AbsolutePath} failed because no ApiCredentials were provided");
|
log.Write(LogLevel.Warning, $"[{requestId}] Request {uri.AbsolutePath} failed because no ApiCredentials were provided");
|
||||||
return new WebCallResult<T>(new NoApiCredentialsError());
|
return new CallResult<IRequest>(new NoApiCredentialsError());
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Write(LogLevel.Information, $"[{requestId}] Creating request for " + uri);
|
log.Write(LogLevel.Information, $"[{requestId}] Creating request for " + uri);
|
||||||
var paramsPosition = parameterPosition ?? apiClient.ParameterPositions[method];
|
var paramsPosition = parameterPosition ?? apiClient.ParameterPositions[method];
|
||||||
var request = ConstructRequest(apiClient, uri, method, parameters, signed, paramsPosition, arraySerialization ?? apiClient.arraySerialization, requestId, additionalHeaders);
|
var request = ConstructRequest(apiClient, uri, method, parameters, signed, paramsPosition, arraySerialization ?? apiClient.arraySerialization, requestId, additionalHeaders);
|
||||||
|
|
||||||
string? paramString = "";
|
string? paramString = "";
|
||||||
if (paramsPosition == HttpMethodParameterPosition.InBody)
|
if (paramsPosition == HttpMethodParameterPosition.InBody)
|
||||||
paramString = $" with request body '{request.Content}'";
|
paramString = $" with request body '{request.Content}'";
|
||||||
@@ -133,12 +213,14 @@ namespace CryptoExchange.Net
|
|||||||
var headers = request.GetHeaders();
|
var headers = request.GetHeaders();
|
||||||
if (headers.Any())
|
if (headers.Any())
|
||||||
paramString += " with headers " + string.Join(", ", headers.Select(h => h.Key + $"=[{string.Join(",", h.Value)}]"));
|
paramString += " with headers " + string.Join(", ", headers.Select(h => h.Key + $"=[{string.Join(",", h.Value)}]"));
|
||||||
|
|
||||||
apiClient.TotalRequestsMade++;
|
apiClient.TotalRequestsMade++;
|
||||||
log.Write(LogLevel.Trace, $"[{requestId}] Sending {method}{(signed ? " signed" : "")} request to {request.Uri}{paramString ?? " "}{(ClientOptions.Proxy == null ? "" : $" via proxy {ClientOptions.Proxy.Host}")}");
|
log.Write(LogLevel.Trace, $"[{requestId}] Sending {method}{(signed ? " signed" : "")} request to {request.Uri}{paramString ?? " "}{(ClientOptions.Proxy == null ? "" : $" via proxy {ClientOptions.Proxy.Host}")}");
|
||||||
return await GetResponseAsync<T>(apiClient, request, deserializer, cancellationToken).ConfigureAwait(false);
|
return new CallResult<IRequest>(request);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Executes the request and returns the result deserialized into the type parameter class
|
/// Executes the request and returns the result deserialized into the type parameter class
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -146,8 +228,14 @@ namespace CryptoExchange.Net
|
|||||||
/// <param name="request">The request object to execute</param>
|
/// <param name="request">The request object to execute</param>
|
||||||
/// <param name="deserializer">The JsonSerializer to use for deserialization</param>
|
/// <param name="deserializer">The JsonSerializer to use for deserialization</param>
|
||||||
/// <param name="cancellationToken">Cancellation token</param>
|
/// <param name="cancellationToken">Cancellation token</param>
|
||||||
|
/// <param name="expectedEmptyResponse">If an empty response is expected</param>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected virtual async Task<WebCallResult<T>> GetResponseAsync<T>(BaseApiClient apiClient, IRequest request, JsonSerializer? deserializer, CancellationToken cancellationToken)
|
protected virtual async Task<WebCallResult<T>> GetResponseAsync<T>(
|
||||||
|
BaseApiClient apiClient,
|
||||||
|
IRequest request,
|
||||||
|
JsonSerializer? deserializer,
|
||||||
|
CancellationToken cancellationToken,
|
||||||
|
bool expectedEmptyResponse)
|
||||||
{
|
{
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
@@ -169,22 +257,52 @@ namespace CryptoExchange.Net
|
|||||||
response.Close();
|
response.Close();
|
||||||
log.Write(LogLevel.Debug, $"[{request.RequestId}] Response received in {sw.ElapsedMilliseconds}ms{(log.Level == LogLevel.Trace ? (": "+data): "")}");
|
log.Write(LogLevel.Debug, $"[{request.RequestId}] Response received in {sw.ElapsedMilliseconds}ms{(log.Level == LogLevel.Trace ? (": "+data): "")}");
|
||||||
|
|
||||||
// Validate if it is valid json. Sometimes other data will be returned, 502 error html pages for example
|
if (!expectedEmptyResponse)
|
||||||
var parseResult = ValidateJson(data);
|
{
|
||||||
if (!parseResult.Success)
|
// Validate if it is valid json. Sometimes other data will be returned, 502 error html pages for example
|
||||||
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, parseResult.Error!);
|
var parseResult = ValidateJson(data);
|
||||||
|
if (!parseResult.Success)
|
||||||
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, parseResult.Error!);
|
||||||
|
|
||||||
// Let the library implementation see if it is an error response, and if so parse the error
|
// Let the library implementation see if it is an error response, and if so parse the error
|
||||||
var error = await TryParseErrorAsync(parseResult.Data).ConfigureAwait(false);
|
var error = await TryParseErrorAsync(parseResult.Data).ConfigureAwait(false);
|
||||||
if (error != null)
|
if (error != null)
|
||||||
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, error!);
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, error!);
|
||||||
|
|
||||||
// Not an error, so continue deserializing
|
// Not an error, so continue deserializing
|
||||||
var deserializeResult = Deserialize<T>(parseResult.Data, deserializer, request.RequestId);
|
var deserializeResult = Deserialize<T>(parseResult.Data, deserializer, request.RequestId);
|
||||||
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data: null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), deserializeResult.Data, deserializeResult.Error);
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), deserializeResult.Data, deserializeResult.Error);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
if (!string.IsNullOrEmpty(data))
|
||||||
|
{
|
||||||
|
var parseResult = ValidateJson(data);
|
||||||
|
if (!parseResult.Success)
|
||||||
|
// Not empty, and not json
|
||||||
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, parseResult.Error!);
|
||||||
|
|
||||||
|
var error = await TryParseErrorAsync(parseResult.Data).ConfigureAwait(false);
|
||||||
|
if (error != null)
|
||||||
|
// Error response
|
||||||
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, error!);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Empty success response; okay
|
||||||
|
return new WebCallResult<T>(response.StatusCode, response.ResponseHeaders, sw.Elapsed, ClientOptions.OutputOriginalData ? data : null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, default);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
|
if (expectedEmptyResponse)
|
||||||
|
{
|
||||||
|
// We expected an empty response and the request is successful and don't manually parse errors, so assume it's correct
|
||||||
|
responseStream.Close();
|
||||||
|
response.Close();
|
||||||
|
|
||||||
|
return new WebCallResult<T>(statusCode, headers, sw.Elapsed, null, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, null);
|
||||||
|
}
|
||||||
|
|
||||||
// Success status code, and we don't have to check for errors. Continue deserializing directly from the stream
|
// Success status code, and we don't have to check for errors. Continue deserializing directly from the stream
|
||||||
var desResult = await DeserializeAsync<T>(responseStream, deserializer, request.RequestId, sw.ElapsedMilliseconds).ConfigureAwait(false);
|
var desResult = await DeserializeAsync<T>(responseStream, deserializer, request.RequestId, sw.ElapsedMilliseconds).ConfigureAwait(false);
|
||||||
responseStream.Close();
|
responseStream.Close();
|
||||||
|
|||||||
@@ -30,15 +30,15 @@ namespace CryptoExchange.Net
|
|||||||
/// <summary>
|
/// <summary>
|
||||||
/// List of socket connections currently connecting/connected
|
/// List of socket connections currently connecting/connected
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected internal ConcurrentDictionary<int, SocketConnection> sockets = new();
|
protected internal ConcurrentDictionary<int, SocketConnection> socketConnections = new();
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Semaphore used while creating sockets
|
/// Semaphore used while creating sockets
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected internal readonly SemaphoreSlim semaphoreSlim = new(1);
|
protected internal readonly SemaphoreSlim semaphoreSlim = new(1);
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The max amount of concurrent socket connections
|
/// Keep alive interval for websocket connection
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected int MaxSocketConnections { get; set; } = 9999;
|
protected TimeSpan KeepAliveInterval { get; set; } = TimeSpan.FromSeconds(10);
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Delegate used for processing byte data received from socket connections before it is processed by handlers
|
/// Delegate used for processing byte data received from socket connections before it is processed by handlers
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -85,10 +85,24 @@ namespace CryptoExchange.Net
|
|||||||
{
|
{
|
||||||
get
|
get
|
||||||
{
|
{
|
||||||
if (!sockets.Any())
|
if (!socketConnections.Any())
|
||||||
return 0;
|
return 0;
|
||||||
|
|
||||||
return sockets.Sum(s => s.Value.Socket.IncomingKbps);
|
return socketConnections.Sum(s => s.Value.IncomingKbps);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public int CurrentConnections => socketConnections.Count;
|
||||||
|
/// <inheritdoc />
|
||||||
|
public int CurrentSubscriptions
|
||||||
|
{
|
||||||
|
get
|
||||||
|
{
|
||||||
|
if (!socketConnections.Any())
|
||||||
|
return 0;
|
||||||
|
|
||||||
|
return socketConnections.Sum(s => s.Value.SubscriptionCount);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -164,7 +178,7 @@ namespace CryptoExchange.Net
|
|||||||
return new CallResult<UpdateSubscription>(new InvalidOperationError("Client disposed, can't subscribe"));
|
return new CallResult<UpdateSubscription>(new InvalidOperationError("Client disposed, can't subscribe"));
|
||||||
|
|
||||||
SocketConnection socketConnection;
|
SocketConnection socketConnection;
|
||||||
SocketSubscription subscription;
|
SocketSubscription? subscription;
|
||||||
var released = false;
|
var released = false;
|
||||||
// Wait for a semaphore here, so we only connect 1 socket at a time.
|
// Wait for a semaphore here, so we only connect 1 socket at a time.
|
||||||
// This is necessary for being able to see if connections can be combined
|
// This is necessary for being able to see if connections can be combined
|
||||||
@@ -179,23 +193,34 @@ namespace CryptoExchange.Net
|
|||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
// Get a new or existing socket connection
|
while (true)
|
||||||
socketConnection = GetSocketConnection(apiClient, url, authenticated);
|
|
||||||
|
|
||||||
// Add a subscription on the socket connection
|
|
||||||
subscription = AddSubscription(request, identifier, true, socketConnection, dataHandler);
|
|
||||||
if (ClientOptions.SocketSubscriptionsCombineTarget == 1)
|
|
||||||
{
|
{
|
||||||
// Only 1 subscription per connection, so no need to wait for connection since a new subscription will create a new connection anyway
|
// Get a new or existing socket connection
|
||||||
semaphoreSlim.Release();
|
socketConnection = GetSocketConnection(apiClient, url, authenticated);
|
||||||
released = true;
|
|
||||||
|
// Add a subscription on the socket connection
|
||||||
|
subscription = AddSubscription(request, identifier, true, socketConnection, dataHandler);
|
||||||
|
if (subscription == null)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {socketConnection.SocketId} failed to add subscription, retrying on different connection");
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (ClientOptions.SocketSubscriptionsCombineTarget == 1)
|
||||||
|
{
|
||||||
|
// Only 1 subscription per connection, so no need to wait for connection since a new subscription will create a new connection anyway
|
||||||
|
semaphoreSlim.Release();
|
||||||
|
released = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
var needsConnecting = !socketConnection.Connected;
|
||||||
|
|
||||||
|
var connectResult = await ConnectIfNeededAsync(socketConnection, authenticated).ConfigureAwait(false);
|
||||||
|
if (!connectResult)
|
||||||
|
return new CallResult<UpdateSubscription>(connectResult.Error!);
|
||||||
|
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
var needsConnecting = !socketConnection.Connected;
|
|
||||||
|
|
||||||
var connectResult = await ConnectIfNeededAsync(socketConnection, authenticated).ConfigureAwait(false);
|
|
||||||
if (!connectResult)
|
|
||||||
return new CallResult<UpdateSubscription>(connectResult.Error!);
|
|
||||||
}
|
}
|
||||||
finally
|
finally
|
||||||
{
|
{
|
||||||
@@ -205,7 +230,7 @@ namespace CryptoExchange.Net
|
|||||||
|
|
||||||
if (socketConnection.PausedActivity)
|
if (socketConnection.PausedActivity)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"Socket {socketConnection.Socket.Id} has been paused, can't subscribe at this moment");
|
log.Write(LogLevel.Warning, $"Socket {socketConnection.SocketId} has been paused, can't subscribe at this moment");
|
||||||
return new CallResult<UpdateSubscription>( new ServerError("Socket is paused"));
|
return new CallResult<UpdateSubscription>( new ServerError("Socket is paused"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -230,12 +255,12 @@ namespace CryptoExchange.Net
|
|||||||
{
|
{
|
||||||
subscription.CancellationTokenRegistration = ct.Register(async () =>
|
subscription.CancellationTokenRegistration = ct.Register(async () =>
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Information, $"Socket {socketConnection.Socket.Id} Cancellation token set, closing subscription");
|
log.Write(LogLevel.Information, $"Socket {socketConnection.SocketId} Cancellation token set, closing subscription");
|
||||||
await socketConnection.CloseAsync(subscription).ConfigureAwait(false);
|
await socketConnection.CloseAsync(subscription).ConfigureAwait(false);
|
||||||
}, false);
|
}, false);
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Write(LogLevel.Information, $"Socket {socketConnection.Socket.Id} subscription completed");
|
log.Write(LogLevel.Information, $"Socket {socketConnection.SocketId} subscription completed");
|
||||||
return new CallResult<UpdateSubscription>(new UpdateSubscription(socketConnection, subscription));
|
return new CallResult<UpdateSubscription>(new UpdateSubscription(socketConnection, subscription));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -317,7 +342,7 @@ namespace CryptoExchange.Net
|
|||||||
|
|
||||||
if (socketConnection.PausedActivity)
|
if (socketConnection.PausedActivity)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"Socket {socketConnection.Socket.Id} has been paused, can't send query at this moment");
|
log.Write(LogLevel.Warning, $"Socket {socketConnection.SocketId} has been paused, can't send query at this moment");
|
||||||
return new CallResult<T>(new ServerError("Socket is paused"));
|
return new CallResult<T>(new ServerError("Socket is paused"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -368,7 +393,7 @@ namespace CryptoExchange.Net
|
|||||||
if (!result)
|
if (!result)
|
||||||
{
|
{
|
||||||
await socket.CloseAsync().ConfigureAwait(false);
|
await socket.CloseAsync().ConfigureAwait(false);
|
||||||
log.Write(LogLevel.Warning, $"Socket {socket.Socket.Id} authentication failed");
|
log.Write(LogLevel.Warning, $"Socket {socket.SocketId} authentication failed");
|
||||||
result.Error!.Message = "Authentication failed: " + result.Error.Message;
|
result.Error!.Message = "Authentication failed: " + result.Error.Message;
|
||||||
return new CallResult<bool>(result.Error);
|
return new CallResult<bool>(result.Error);
|
||||||
}
|
}
|
||||||
@@ -457,7 +482,7 @@ namespace CryptoExchange.Net
|
|||||||
/// <param name="connection">The socket connection the handler is on</param>
|
/// <param name="connection">The socket connection the handler is on</param>
|
||||||
/// <param name="dataHandler">The handler of the data received</param>
|
/// <param name="dataHandler">The handler of the data received</param>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected virtual SocketSubscription AddSubscription<T>(object? request, string? identifier, bool userSubscription, SocketConnection connection, Action<DataEvent<T>> dataHandler)
|
protected virtual SocketSubscription? AddSubscription<T>(object? request, string? identifier, bool userSubscription, SocketConnection connection, Action<DataEvent<T>> dataHandler)
|
||||||
{
|
{
|
||||||
void InternalHandler(MessageEvent messageEvent)
|
void InternalHandler(MessageEvent messageEvent)
|
||||||
{
|
{
|
||||||
@@ -471,7 +496,7 @@ namespace CryptoExchange.Net
|
|||||||
var desResult = Deserialize<T>(messageEvent.JsonData);
|
var desResult = Deserialize<T>(messageEvent.JsonData);
|
||||||
if (!desResult)
|
if (!desResult)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"Socket {connection.Socket.Id} Failed to deserialize data into type {typeof(T)}: {desResult.Error}");
|
log.Write(LogLevel.Warning, $"Socket {connection.SocketId} Failed to deserialize data into type {typeof(T)}: {desResult.Error}");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -481,7 +506,8 @@ namespace CryptoExchange.Net
|
|||||||
var subscription = request == null
|
var subscription = request == null
|
||||||
? SocketSubscription.CreateForIdentifier(NextId(), identifier!, userSubscription, InternalHandler)
|
? SocketSubscription.CreateForIdentifier(NextId(), identifier!, userSubscription, InternalHandler)
|
||||||
: SocketSubscription.CreateForRequest(NextId(), request, userSubscription, InternalHandler);
|
: SocketSubscription.CreateForRequest(NextId(), request, userSubscription, InternalHandler);
|
||||||
connection.AddSubscription(subscription);
|
if (!connection.AddSubscription(subscription))
|
||||||
|
return null;
|
||||||
return subscription;
|
return subscription;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -494,7 +520,7 @@ namespace CryptoExchange.Net
|
|||||||
{
|
{
|
||||||
genericHandlers.Add(identifier, action);
|
genericHandlers.Add(identifier, action);
|
||||||
var subscription = SocketSubscription.CreateForIdentifier(NextId(), identifier, false, action);
|
var subscription = SocketSubscription.CreateForIdentifier(NextId(), identifier, false, action);
|
||||||
foreach (var connection in sockets.Values)
|
foreach (var connection in socketConnections.Values)
|
||||||
connection.AddSubscription(subscription);
|
connection.AddSubscription(subscription);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -507,13 +533,14 @@ namespace CryptoExchange.Net
|
|||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected virtual SocketConnection GetSocketConnection(SocketApiClient apiClient, string address, bool authenticated)
|
protected virtual SocketConnection GetSocketConnection(SocketApiClient apiClient, string address, bool authenticated)
|
||||||
{
|
{
|
||||||
var socketResult = sockets.Where(s => s.Value.Socket.Url.TrimEnd('/') == address.TrimEnd('/')
|
var socketResult = socketConnections.Where(s => (s.Value.Status == SocketConnection.SocketStatus.None || s.Value.Status == SocketConnection.SocketStatus.Connected)
|
||||||
|
&& s.Value.Uri.ToString().TrimEnd('/') == address.TrimEnd('/')
|
||||||
&& (s.Value.ApiClient.GetType() == apiClient.GetType())
|
&& (s.Value.ApiClient.GetType() == apiClient.GetType())
|
||||||
&& (s.Value.Authenticated == authenticated || !authenticated) && s.Value.Connected).OrderBy(s => s.Value.SubscriptionCount).FirstOrDefault();
|
&& (s.Value.Authenticated == authenticated || !authenticated) && s.Value.Connected).OrderBy(s => s.Value.SubscriptionCount).FirstOrDefault();
|
||||||
var result = socketResult.Equals(default(KeyValuePair<int, SocketConnection>)) ? null : socketResult.Value;
|
var result = socketResult.Equals(default(KeyValuePair<int, SocketConnection>)) ? null : socketResult.Value;
|
||||||
if (result != null)
|
if (result != null)
|
||||||
{
|
{
|
||||||
if (result.SubscriptionCount < ClientOptions.SocketSubscriptionsCombineTarget || (sockets.Count >= MaxSocketConnections && sockets.All(s => s.Value.SubscriptionCount >= ClientOptions.SocketSubscriptionsCombineTarget)))
|
if (result.SubscriptionCount < ClientOptions.SocketSubscriptionsCombineTarget || (socketConnections.Count >= ClientOptions.MaxSocketConnections && socketConnections.All(s => s.Value.SubscriptionCount >= ClientOptions.SocketSubscriptionsCombineTarget)))
|
||||||
{
|
{
|
||||||
// Use existing socket if it has less than target connections OR it has the least connections and we can't make new
|
// Use existing socket if it has less than target connections OR it has the least connections and we can't make new
|
||||||
return result;
|
return result;
|
||||||
@@ -548,13 +575,13 @@ namespace CryptoExchange.Net
|
|||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected virtual async Task<CallResult<bool>> ConnectSocketAsync(SocketConnection socketConnection)
|
protected virtual async Task<CallResult<bool>> ConnectSocketAsync(SocketConnection socketConnection)
|
||||||
{
|
{
|
||||||
if (await socketConnection.Socket.ConnectAsync().ConfigureAwait(false))
|
if (await socketConnection.ConnectAsync().ConfigureAwait(false))
|
||||||
{
|
{
|
||||||
sockets.TryAdd(socketConnection.Socket.Id, socketConnection);
|
socketConnections.TryAdd(socketConnection.SocketId, socketConnection);
|
||||||
return new CallResult<bool>(true);
|
return new CallResult<bool>(true);
|
||||||
}
|
}
|
||||||
|
|
||||||
socketConnection.Socket.Dispose();
|
socketConnection.Dispose();
|
||||||
return new CallResult<bool>(new CantConnectError());
|
return new CallResult<bool>(new CantConnectError());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -571,6 +598,7 @@ namespace CryptoExchange.Net
|
|||||||
if (ClientOptions.Proxy != null)
|
if (ClientOptions.Proxy != null)
|
||||||
socket.SetProxy(ClientOptions.Proxy);
|
socket.SetProxy(ClientOptions.Proxy);
|
||||||
|
|
||||||
|
socket.KeepAliveInterval = KeepAliveInterval;
|
||||||
socket.Timeout = ClientOptions.SocketNoDataTimeout;
|
socket.Timeout = ClientOptions.SocketNoDataTimeout;
|
||||||
socket.DataInterpreterBytes = dataInterpreterBytes;
|
socket.DataInterpreterBytes = dataInterpreterBytes;
|
||||||
socket.DataInterpreterString = dataInterpreterString;
|
socket.DataInterpreterString = dataInterpreterString;
|
||||||
@@ -605,27 +633,27 @@ namespace CryptoExchange.Net
|
|||||||
if (disposing)
|
if (disposing)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
foreach (var socket in sockets.Values)
|
foreach (var socketConnection in socketConnections.Values)
|
||||||
{
|
{
|
||||||
if (disposing)
|
if (disposing)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
if (!socket.Socket.IsOpen)
|
if (!socketConnection.Connected)
|
||||||
continue;
|
continue;
|
||||||
|
|
||||||
var obj = objGetter(socket);
|
var obj = objGetter(socketConnection);
|
||||||
if (obj == null)
|
if (obj == null)
|
||||||
continue;
|
continue;
|
||||||
|
|
||||||
log.Write(LogLevel.Trace, $"Socket {socket.Socket.Id} sending periodic {identifier}");
|
log.Write(LogLevel.Trace, $"Socket {socketConnection.SocketId} sending periodic {identifier}");
|
||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
socket.Send(obj);
|
socketConnection.Send(obj);
|
||||||
}
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"Socket {socket.Socket.Id} Periodic send {identifier} failed: " + ex.ToLogString());
|
log.Write(LogLevel.Warning, $"Socket {socketConnection.SocketId} Periodic send {identifier} failed: " + ex.ToLogString());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -642,7 +670,7 @@ namespace CryptoExchange.Net
|
|||||||
|
|
||||||
SocketSubscription? subscription = null;
|
SocketSubscription? subscription = null;
|
||||||
SocketConnection? connection = null;
|
SocketConnection? connection = null;
|
||||||
foreach(var socket in sockets.Values.ToList())
|
foreach(var socket in socketConnections.Values.ToList())
|
||||||
{
|
{
|
||||||
subscription = socket.GetSubscription(subscriptionId);
|
subscription = socket.GetSubscription(subscriptionId);
|
||||||
if (subscription != null)
|
if (subscription != null)
|
||||||
@@ -679,19 +707,15 @@ namespace CryptoExchange.Net
|
|||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
public virtual async Task UnsubscribeAllAsync()
|
public virtual async Task UnsubscribeAllAsync()
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Information, $"Closing all {sockets.Sum(s => s.Value.SubscriptionCount)} subscriptions");
|
log.Write(LogLevel.Information, $"Closing all {socketConnections.Sum(s => s.Value.SubscriptionCount)} subscriptions");
|
||||||
|
var tasks = new List<Task>();
|
||||||
await Task.Run(async () =>
|
|
||||||
{
|
{
|
||||||
var tasks = new List<Task>();
|
var socketList = socketConnections.Values;
|
||||||
{
|
foreach (var sub in socketList)
|
||||||
var socketList = sockets.Values;
|
tasks.Add(sub.CloseAsync());
|
||||||
foreach (var sub in socketList)
|
}
|
||||||
tasks.Add(sub.CloseAsync());
|
|
||||||
}
|
|
||||||
|
|
||||||
await Task.WhenAll(tasks.ToArray()).ConfigureAwait(false);
|
await Task.WhenAll(tasks.ToArray()).ConfigureAwait(false);
|
||||||
}).ConfigureAwait(false);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -703,7 +727,7 @@ namespace CryptoExchange.Net
|
|||||||
periodicEvent?.Set();
|
periodicEvent?.Set();
|
||||||
periodicEvent?.Dispose();
|
periodicEvent?.Dispose();
|
||||||
log.Write(LogLevel.Debug, "Disposing socket client, closing all subscriptions");
|
log.Write(LogLevel.Debug, "Disposing socket client, closing all subscriptions");
|
||||||
Task.Run(UnsubscribeAllAsync).ConfigureAwait(false).GetAwaiter().GetResult();
|
_ = UnsubscribeAllAsync();
|
||||||
semaphoreSlim?.Dispose();
|
semaphoreSlim?.Dispose();
|
||||||
base.Dispose();
|
base.Dispose();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,7 +18,7 @@ namespace CryptoExchange.Net
|
|||||||
/// Get time sync info for an API client
|
/// Get time sync info for an API client
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected abstract TimeSyncInfo GetTimeSyncInfo();
|
public abstract TimeSyncInfo GetTimeSyncInfo();
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Get time offset for an API client
|
/// Get time offset for an API client
|
||||||
@@ -92,7 +92,7 @@ namespace CryptoExchange.Net
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Calculate time offset between local and server
|
// Calculate time offset between local and server
|
||||||
var offset = result.Data - localTime;
|
var offset = result.Data - (localTime.AddMilliseconds(result.ResponseTime!.Value.TotalMilliseconds / 2));
|
||||||
timeSyncParams.UpdateTimeOffset(offset);
|
timeSyncParams.UpdateTimeOffset(offset);
|
||||||
timeSyncParams.TimeSyncState.Semaphore.Release();
|
timeSyncParams.TimeSyncState.Semaphore.Release();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
using Newtonsoft.Json;
|
using Newtonsoft.Json;
|
||||||
using Newtonsoft.Json.Linq;
|
|
||||||
using System;
|
using System;
|
||||||
using System.Diagnostics;
|
using System.Diagnostics;
|
||||||
using System.Diagnostics.CodeAnalysis;
|
using System.Diagnostics.CodeAnalysis;
|
||||||
@@ -32,13 +31,13 @@ namespace CryptoExchange.Net.Converters
|
|||||||
if(reader.TokenType is JsonToken.Integer)
|
if(reader.TokenType is JsonToken.Integer)
|
||||||
{
|
{
|
||||||
var longValue = (long)reader.Value;
|
var longValue = (long)reader.Value;
|
||||||
if (longValue == 0)
|
if (longValue == 0 || longValue == -1)
|
||||||
return objectType == typeof(DateTime) ? default(DateTime): null;
|
return objectType == typeof(DateTime) ? default(DateTime): null;
|
||||||
if (longValue < 1999999999)
|
if (longValue < 19999999999)
|
||||||
return ConvertFromSeconds(longValue);
|
return ConvertFromSeconds(longValue);
|
||||||
if (longValue < 1999999999999)
|
if (longValue < 19999999999999)
|
||||||
return ConvertFromMilliseconds(longValue);
|
return ConvertFromMilliseconds(longValue);
|
||||||
if (longValue < 1999999999999999)
|
if (longValue < 19999999999999999)
|
||||||
return ConvertFromMicroseconds(longValue);
|
return ConvertFromMicroseconds(longValue);
|
||||||
|
|
||||||
return ConvertFromNanoseconds(longValue);
|
return ConvertFromNanoseconds(longValue);
|
||||||
@@ -46,7 +45,10 @@ namespace CryptoExchange.Net.Converters
|
|||||||
else if (reader.TokenType is JsonToken.Float)
|
else if (reader.TokenType is JsonToken.Float)
|
||||||
{
|
{
|
||||||
var doubleValue = (double)reader.Value;
|
var doubleValue = (double)reader.Value;
|
||||||
if (doubleValue < 1999999999)
|
if (doubleValue == 0 || doubleValue == -1)
|
||||||
|
return objectType == typeof(DateTime) ? default(DateTime) : null;
|
||||||
|
|
||||||
|
if (doubleValue < 19999999999)
|
||||||
return ConvertFromSeconds(doubleValue);
|
return ConvertFromSeconds(doubleValue);
|
||||||
|
|
||||||
return ConvertFromMilliseconds(doubleValue);
|
return ConvertFromMilliseconds(doubleValue);
|
||||||
@@ -57,6 +59,9 @@ namespace CryptoExchange.Net.Converters
|
|||||||
if (string.IsNullOrWhiteSpace(stringValue))
|
if (string.IsNullOrWhiteSpace(stringValue))
|
||||||
return null;
|
return null;
|
||||||
|
|
||||||
|
if (string.IsNullOrWhiteSpace(stringValue) || stringValue == "0" || stringValue == "-1")
|
||||||
|
return objectType == typeof(DateTime) ? default(DateTime) : null;
|
||||||
|
|
||||||
if (stringValue.Length == 8)
|
if (stringValue.Length == 8)
|
||||||
{
|
{
|
||||||
// Parse 20211103 format
|
// Parse 20211103 format
|
||||||
@@ -86,11 +91,11 @@ namespace CryptoExchange.Net.Converters
|
|||||||
if (double.TryParse(stringValue, NumberStyles.Float, CultureInfo.InvariantCulture, out var doubleValue))
|
if (double.TryParse(stringValue, NumberStyles.Float, CultureInfo.InvariantCulture, out var doubleValue))
|
||||||
{
|
{
|
||||||
// Parse 1637745563.000 format
|
// Parse 1637745563.000 format
|
||||||
if (doubleValue < 1999999999)
|
if (doubleValue < 19999999999)
|
||||||
return ConvertFromSeconds(doubleValue);
|
return ConvertFromSeconds(doubleValue);
|
||||||
if (doubleValue < 1999999999999)
|
if (doubleValue < 19999999999999)
|
||||||
return ConvertFromMilliseconds((long)doubleValue);
|
return ConvertFromMilliseconds((long)doubleValue);
|
||||||
if (doubleValue < 1999999999999999)
|
if (doubleValue < 19999999999999999)
|
||||||
return ConvertFromMicroseconds((long)doubleValue);
|
return ConvertFromMicroseconds((long)doubleValue);
|
||||||
|
|
||||||
return ConvertFromNanoseconds((long)doubleValue);
|
return ConvertFromNanoseconds((long)doubleValue);
|
||||||
|
|||||||
@@ -6,16 +6,16 @@
|
|||||||
<PackageId>CryptoExchange.Net</PackageId>
|
<PackageId>CryptoExchange.Net</PackageId>
|
||||||
<Authors>JKorf</Authors>
|
<Authors>JKorf</Authors>
|
||||||
<Description>A base package for implementing cryptocurrency API's</Description>
|
<Description>A base package for implementing cryptocurrency API's</Description>
|
||||||
<PackageVersion>5.1.7</PackageVersion>
|
<PackageVersion>5.1.12</PackageVersion>
|
||||||
<AssemblyVersion>5.1.7</AssemblyVersion>
|
<AssemblyVersion>5.1.12</AssemblyVersion>
|
||||||
<FileVersion>5.1.7</FileVersion>
|
<FileVersion>5.1.12</FileVersion>
|
||||||
<PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance>
|
<PackageRequireLicenseAcceptance>false</PackageRequireLicenseAcceptance>
|
||||||
<RepositoryType>git</RepositoryType>
|
<RepositoryType>git</RepositoryType>
|
||||||
<RepositoryUrl>https://github.com/JKorf/CryptoExchange.Net.git</RepositoryUrl>
|
<RepositoryUrl>https://github.com/JKorf/CryptoExchange.Net.git</RepositoryUrl>
|
||||||
<PackageProjectUrl>https://github.com/JKorf/CryptoExchange.Net</PackageProjectUrl>
|
<PackageProjectUrl>https://github.com/JKorf/CryptoExchange.Net</PackageProjectUrl>
|
||||||
<NeutralLanguage>en</NeutralLanguage>
|
<NeutralLanguage>en</NeutralLanguage>
|
||||||
<GeneratePackageOnBuild>true</GeneratePackageOnBuild>
|
<GeneratePackageOnBuild>true</GeneratePackageOnBuild>
|
||||||
<PackageReleaseNotes>5.1.7 - Moved some Rest parameters from BaseRestClient to RestApiClient to allow different implementations for sub clients</PackageReleaseNotes>
|
<PackageReleaseNotes>5.1.12 - Changed time sync so requests no longer wait for it to complete unless it's the first time, Made log client options changable after client creation, Fixed proxy setting not used when reconnecting socket, Changed MaxSocketConnections to a client options, Updated socket reconnection logic</PackageReleaseNotes>
|
||||||
<Nullable>enable</Nullable>
|
<Nullable>enable</Nullable>
|
||||||
<LangVersion>9.0</LangVersion>
|
<LangVersion>9.0</LangVersion>
|
||||||
<PackageLicenseExpression>MIT</PackageLicenseExpression>
|
<PackageLicenseExpression>MIT</PackageLicenseExpression>
|
||||||
|
|||||||
@@ -426,6 +426,7 @@ namespace CryptoExchange.Net
|
|||||||
var uriBuilder = new UriBuilder();
|
var uriBuilder = new UriBuilder();
|
||||||
uriBuilder.Scheme = baseUri.Scheme;
|
uriBuilder.Scheme = baseUri.Scheme;
|
||||||
uriBuilder.Host = baseUri.Host;
|
uriBuilder.Host = baseUri.Host;
|
||||||
|
uriBuilder.Port = baseUri.Port;
|
||||||
uriBuilder.Path = baseUri.AbsolutePath;
|
uriBuilder.Path = baseUri.AbsolutePath;
|
||||||
var httpValueCollection = HttpUtility.ParseQueryString(string.Empty);
|
var httpValueCollection = HttpUtility.ParseQueryString(string.Empty);
|
||||||
foreach (var parameter in parameters)
|
foreach (var parameter in parameters)
|
||||||
@@ -454,6 +455,7 @@ namespace CryptoExchange.Net
|
|||||||
var uriBuilder = new UriBuilder();
|
var uriBuilder = new UriBuilder();
|
||||||
uriBuilder.Scheme = baseUri.Scheme;
|
uriBuilder.Scheme = baseUri.Scheme;
|
||||||
uriBuilder.Host = baseUri.Host;
|
uriBuilder.Host = baseUri.Host;
|
||||||
|
uriBuilder.Port = baseUri.Port;
|
||||||
uriBuilder.Path = baseUri.AbsolutePath;
|
uriBuilder.Path = baseUri.AbsolutePath;
|
||||||
var httpValueCollection = HttpUtility.ParseQueryString(string.Empty);
|
var httpValueCollection = HttpUtility.ParseQueryString(string.Empty);
|
||||||
foreach (var parameter in parameters)
|
foreach (var parameter in parameters)
|
||||||
|
|||||||
@@ -27,6 +27,16 @@ namespace CryptoExchange.Net.Interfaces
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public double IncomingKbps { get; }
|
public double IncomingKbps { get; }
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The current amount of connections to the API from this client. A connection can have multiple subscriptions.
|
||||||
|
/// </summary>
|
||||||
|
public int CurrentConnections { get; }
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The current amount of subscriptions running from the client
|
||||||
|
/// </summary>
|
||||||
|
public int CurrentSubscriptions { get; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Unsubscribe from a stream using the subscription id received when starting the subscription
|
/// Unsubscribe from a stream using the subscription id received when starting the subscription
|
||||||
/// </summary>
|
/// </summary>
|
||||||
|
|||||||
@@ -41,10 +41,6 @@ namespace CryptoExchange.Net.Interfaces
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
Encoding? Encoding { get; set; }
|
Encoding? Encoding { get; set; }
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Whether socket is in the process of reconnecting
|
|
||||||
/// </summary>
|
|
||||||
bool Reconnecting { get; set; }
|
|
||||||
/// <summary>
|
|
||||||
/// The max amount of outgoing messages per second
|
/// The max amount of outgoing messages per second
|
||||||
/// </summary>
|
/// </summary>
|
||||||
int? RatelimitPerSecond { get; set; }
|
int? RatelimitPerSecond { get; set; }
|
||||||
@@ -61,9 +57,9 @@ namespace CryptoExchange.Net.Interfaces
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
Func<string, string>? DataInterpreterString { get; set; }
|
Func<string, string>? DataInterpreterString { get; set; }
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The url the socket connects to
|
/// The uri the socket connects to
|
||||||
/// </summary>
|
/// </summary>
|
||||||
string Url { get; }
|
Uri Uri { get; }
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Whether the socket connection is closed
|
/// Whether the socket connection is closed
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -81,6 +77,10 @@ namespace CryptoExchange.Net.Interfaces
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
TimeSpan Timeout { get; set; }
|
TimeSpan Timeout { get; set; }
|
||||||
/// <summary>
|
/// <summary>
|
||||||
|
/// The interval at which to send a ping frame to the server
|
||||||
|
/// </summary>
|
||||||
|
TimeSpan KeepAliveInterval { get; set; }
|
||||||
|
/// <summary>
|
||||||
/// Set a proxy to use when connecting
|
/// Set a proxy to use when connecting
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="proxy"></param>
|
/// <param name="proxy"></param>
|
||||||
@@ -89,7 +89,12 @@ namespace CryptoExchange.Net.Interfaces
|
|||||||
/// Connect the socket
|
/// Connect the socket
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
Task<bool> ConnectAsync();
|
Task<bool> ConnectAsync();
|
||||||
|
/// <summary>
|
||||||
|
/// Receive and send messages over the connection. Resulting task should complete when closing the socket.
|
||||||
|
/// </summary>
|
||||||
|
/// <returns></returns>
|
||||||
|
Task ProcessAsync();
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Send data
|
/// Send data
|
||||||
/// </summary>
|
/// </summary>
|
||||||
|
|||||||
@@ -26,6 +26,8 @@ namespace CryptoExchange.Net.Logging
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public string ClientName { get; set; }
|
public string ClientName { get; set; }
|
||||||
|
|
||||||
|
private readonly object _lock = new object();
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// ctor
|
/// ctor
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -42,7 +44,8 @@ namespace CryptoExchange.Net.Logging
|
|||||||
/// <param name="textWriters"></param>
|
/// <param name="textWriters"></param>
|
||||||
public void UpdateWriters(List<ILogger> textWriters)
|
public void UpdateWriters(List<ILogger> textWriters)
|
||||||
{
|
{
|
||||||
writers = textWriters;
|
lock (_lock)
|
||||||
|
writers = textWriters;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -56,16 +59,19 @@ namespace CryptoExchange.Net.Logging
|
|||||||
return;
|
return;
|
||||||
|
|
||||||
var logMessage = $"{ClientName,-10} | {message}";
|
var logMessage = $"{ClientName,-10} | {message}";
|
||||||
foreach (var writer in writers.ToList())
|
lock (_lock)
|
||||||
{
|
{
|
||||||
try
|
foreach (var writer in writers)
|
||||||
{
|
{
|
||||||
writer.Log(logLevel, logMessage);
|
try
|
||||||
}
|
{
|
||||||
catch (Exception e)
|
writer.Log(logLevel, logMessage);
|
||||||
{
|
}
|
||||||
// Can't write to the logging so where else to output..
|
catch (Exception e)
|
||||||
Trace.WriteLine($"{DateTime.Now:yyyy/MM/dd HH:mm:ss:fff} | Warning | Failed to write log to writer {writer.GetType()}: " + e.ToLogString());
|
{
|
||||||
|
// Can't write to the logging so where else to output..
|
||||||
|
Trace.WriteLine($"{DateTime.Now:yyyy/MM/dd HH:mm:ss:fff} | Warning | Failed to write log to writer {writer.GetType()}: " + e.ToLogString());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -14,15 +14,35 @@ namespace CryptoExchange.Net.Objects
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public class BaseOptions
|
public class BaseOptions
|
||||||
{
|
{
|
||||||
|
internal event Action? OnLoggingChanged;
|
||||||
|
|
||||||
|
private LogLevel _logLevel = LogLevel.Information;
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The minimum log level to output
|
/// The minimum log level to output
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public LogLevel LogLevel { get; set; } = LogLevel.Information;
|
public LogLevel LogLevel
|
||||||
|
{
|
||||||
|
get => _logLevel;
|
||||||
|
set
|
||||||
|
{
|
||||||
|
_logLevel = value;
|
||||||
|
OnLoggingChanged?.Invoke();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private List<ILogger> _logWriters = new List<ILogger> { new DebugLogger() };
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The log writers
|
/// The log writers
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public List<ILogger> LogWriters { get; set; } = new List<ILogger> { new DebugLogger() };
|
public List<ILogger> LogWriters
|
||||||
|
{
|
||||||
|
get => _logWriters;
|
||||||
|
set
|
||||||
|
{
|
||||||
|
_logWriters = value;
|
||||||
|
OnLoggingChanged?.Invoke();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// If true, the CallResult and DataEvent objects will also include the originally received json data in the OriginalData property
|
/// If true, the CallResult and DataEvent objects will also include the originally received json data in the OriginalData property
|
||||||
@@ -189,6 +209,11 @@ namespace CryptoExchange.Net.Objects
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public int? SocketSubscriptionsCombineTarget { get; set; }
|
public int? SocketSubscriptionsCombineTarget { get; set; }
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The max amount of connections to make to the server. Can be used for API's which only allow a certain number of connections. Changing this to a high value might cause issues.
|
||||||
|
/// </summary>
|
||||||
|
public int? MaxSocketConnections { get; set; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// ctor
|
/// ctor
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -213,12 +238,13 @@ namespace CryptoExchange.Net.Objects
|
|||||||
SocketResponseTimeout = baseOptions.SocketResponseTimeout;
|
SocketResponseTimeout = baseOptions.SocketResponseTimeout;
|
||||||
SocketNoDataTimeout = baseOptions.SocketNoDataTimeout;
|
SocketNoDataTimeout = baseOptions.SocketNoDataTimeout;
|
||||||
SocketSubscriptionsCombineTarget = baseOptions.SocketSubscriptionsCombineTarget;
|
SocketSubscriptionsCombineTarget = baseOptions.SocketSubscriptionsCombineTarget;
|
||||||
|
MaxSocketConnections = baseOptions.MaxSocketConnections;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public override string ToString()
|
public override string ToString()
|
||||||
{
|
{
|
||||||
return $"{base.ToString()}, AutoReconnect: {AutoReconnect}, ReconnectInterval: {ReconnectInterval}, MaxReconnectTries: {MaxReconnectTries}, MaxResubscribeTries: {MaxResubscribeTries}, MaxConcurrentResubscriptionsPerSocket: {MaxConcurrentResubscriptionsPerSocket}, SocketResponseTimeout: {SocketResponseTimeout:c}, SocketNoDataTimeout: {SocketNoDataTimeout}, SocketSubscriptionsCombineTarget: {SocketSubscriptionsCombineTarget}";
|
return $"{base.ToString()}, AutoReconnect: {AutoReconnect}, ReconnectInterval: {ReconnectInterval}, MaxReconnectTries: {MaxReconnectTries}, MaxResubscribeTries: {MaxResubscribeTries}, MaxConcurrentResubscriptionsPerSocket: {MaxConcurrentResubscriptionsPerSocket}, SocketResponseTimeout: {SocketResponseTimeout:c}, SocketNoDataTimeout: {SocketNoDataTimeout}, SocketSubscriptionsCombineTarget: {SocketSubscriptionsCombineTarget}, MaxSocketConnections: {MaxSocketConnections}";
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -258,24 +258,32 @@ namespace CryptoExchange.Net.OrderBook
|
|||||||
}
|
}
|
||||||
|
|
||||||
_subscription = startResult.Data;
|
_subscription = startResult.Data;
|
||||||
_subscription.ConnectionLost += () =>
|
_subscription.ConnectionLost += HandleConnectionLost;
|
||||||
{
|
_subscription.ConnectionClosed += HandleConnectionClosed;
|
||||||
log.Write(LogLevel.Warning, $"{Id} order book {Symbol} connection lost");
|
_subscription.ConnectionRestored += HandleConnectionRestored;
|
||||||
Status = OrderBookStatus.Reconnecting;
|
|
||||||
Reset();
|
|
||||||
};
|
|
||||||
_subscription.ConnectionClosed += () =>
|
|
||||||
{
|
|
||||||
log.Write(LogLevel.Warning, $"{Id} order book {Symbol} disconnected");
|
|
||||||
Status = OrderBookStatus.Disconnected;
|
|
||||||
_ = StopAsync();
|
|
||||||
};
|
|
||||||
|
|
||||||
_subscription.ConnectionRestored += async time => await ResyncAsync().ConfigureAwait(false);
|
|
||||||
Status = OrderBookStatus.Synced;
|
Status = OrderBookStatus.Synced;
|
||||||
return new CallResult<bool>(true);
|
return new CallResult<bool>(true);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void HandleConnectionLost() {
|
||||||
|
log.Write(LogLevel.Warning, $"{Id} order book {Symbol} connection lost");
|
||||||
|
if (Status != OrderBookStatus.Disposed) {
|
||||||
|
Status = OrderBookStatus.Reconnecting;
|
||||||
|
Reset();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void HandleConnectionClosed() {
|
||||||
|
log.Write(LogLevel.Warning, $"{Id} order book {Symbol} disconnected");
|
||||||
|
Status = OrderBookStatus.Disconnected;
|
||||||
|
_ = StopAsync();
|
||||||
|
}
|
||||||
|
|
||||||
|
private async void HandleConnectionRestored(TimeSpan _) {
|
||||||
|
await ResyncAsync().ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
|
||||||
/// <inheritdoc/>
|
/// <inheritdoc/>
|
||||||
public async Task StopAsync()
|
public async Task StopAsync()
|
||||||
{
|
{
|
||||||
@@ -286,8 +294,12 @@ namespace CryptoExchange.Net.OrderBook
|
|||||||
if (_processTask != null)
|
if (_processTask != null)
|
||||||
await _processTask.ConfigureAwait(false);
|
await _processTask.ConfigureAwait(false);
|
||||||
|
|
||||||
if (_subscription != null)
|
if (_subscription != null) {
|
||||||
await _subscription.CloseAsync().ConfigureAwait(false);
|
await _subscription.CloseAsync().ConfigureAwait(false);
|
||||||
|
_subscription.ConnectionLost -= HandleConnectionLost;
|
||||||
|
_subscription.ConnectionClosed -= HandleConnectionClosed;
|
||||||
|
_subscription.ConnectionRestored -= HandleConnectionRestored;
|
||||||
|
}
|
||||||
log.Write(LogLevel.Trace, $"{Id} order book {Symbol} stopped");
|
log.Write(LogLevel.Trace, $"{Id} order book {Symbol} stopped");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -483,7 +495,7 @@ namespace CryptoExchange.Net.OrderBook
|
|||||||
/// <param name="timeout">Max wait time</param>
|
/// <param name="timeout">Max wait time</param>
|
||||||
/// <param name="ct">Cancellation token</param>
|
/// <param name="ct">Cancellation token</param>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
protected async Task<CallResult<bool>> WaitForSetOrderBookAsync(int timeout, CancellationToken ct)
|
protected async Task<CallResult<bool>> WaitForSetOrderBookAsync(TimeSpan timeout, CancellationToken ct)
|
||||||
{
|
{
|
||||||
var startWait = DateTime.UtcNow;
|
var startWait = DateTime.UtcNow;
|
||||||
while (!bookSet && Status == OrderBookStatus.Syncing)
|
while (!bookSet && Status == OrderBookStatus.Syncing)
|
||||||
@@ -491,12 +503,12 @@ namespace CryptoExchange.Net.OrderBook
|
|||||||
if(ct.IsCancellationRequested)
|
if(ct.IsCancellationRequested)
|
||||||
return new CallResult<bool>(new CancellationRequestedError());
|
return new CallResult<bool>(new CancellationRequestedError());
|
||||||
|
|
||||||
if ((DateTime.UtcNow - startWait).TotalMilliseconds > timeout)
|
if (DateTime.UtcNow - startWait > timeout)
|
||||||
return new CallResult<bool>(new ServerError("Timeout while waiting for data"));
|
return new CallResult<bool>(new ServerError("Timeout while waiting for data"));
|
||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
await Task.Delay(10, ct).ConfigureAwait(false);
|
await Task.Delay(50, ct).ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
catch (OperationCanceledException)
|
catch (OperationCanceledException)
|
||||||
{ }
|
{ }
|
||||||
@@ -601,13 +613,13 @@ namespace CryptoExchange.Net.OrderBook
|
|||||||
|
|
||||||
private async Task ProcessQueue()
|
private async Task ProcessQueue()
|
||||||
{
|
{
|
||||||
while (Status != OrderBookStatus.Disconnected)
|
while (Status != OrderBookStatus.Disconnected && Status != OrderBookStatus.Disposed)
|
||||||
{
|
{
|
||||||
await _queueEvent.WaitAsync().ConfigureAwait(false);
|
await _queueEvent.WaitAsync().ConfigureAwait(false);
|
||||||
|
|
||||||
while (_processQueue.TryDequeue(out var item))
|
while (_processQueue.TryDequeue(out var item))
|
||||||
{
|
{
|
||||||
if (Status == OrderBookStatus.Disconnected)
|
if (Status == OrderBookStatus.Disconnected || Status == OrderBookStatus.Disposed)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
if (_stopProcessing)
|
if (_stopProcessing)
|
||||||
|
|||||||
@@ -23,28 +23,26 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
public class CryptoExchangeWebSocketClient : IWebsocket
|
public class CryptoExchangeWebSocketClient : IWebsocket
|
||||||
{
|
{
|
||||||
internal static int lastStreamId;
|
internal static int lastStreamId;
|
||||||
private static readonly object streamIdLock = new object();
|
private static readonly object streamIdLock = new();
|
||||||
|
|
||||||
private ClientWebSocket _socket;
|
private ClientWebSocket _socket;
|
||||||
private Task? _sendTask;
|
|
||||||
private Task? _receiveTask;
|
|
||||||
private Task? _timeoutTask;
|
|
||||||
private readonly AsyncResetEvent _sendEvent;
|
private readonly AsyncResetEvent _sendEvent;
|
||||||
private readonly ConcurrentQueue<byte[]> _sendBuffer;
|
private readonly ConcurrentQueue<byte[]> _sendBuffer;
|
||||||
private readonly IDictionary<string, string> cookies;
|
private readonly IDictionary<string, string> cookies;
|
||||||
private readonly IDictionary<string, string> headers;
|
private readonly IDictionary<string, string> headers;
|
||||||
private CancellationTokenSource _ctsSource;
|
private CancellationTokenSource _ctsSource;
|
||||||
private bool _closing;
|
private ApiProxy? _proxy;
|
||||||
private bool _startedSent;
|
|
||||||
private bool _startedReceive;
|
|
||||||
|
|
||||||
private readonly List<DateTime> _outgoingMessages;
|
private readonly List<DateTime> _outgoingMessages;
|
||||||
private DateTime _lastReceivedMessagesUpdate;
|
private DateTime _lastReceivedMessagesUpdate;
|
||||||
|
private bool _closed;
|
||||||
|
private bool _disposed;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Received messages, the size and the timstamp
|
/// Received messages, the size and the timstamp
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected readonly List<ReceiveItem> _receivedMessages;
|
protected readonly List<ReceiveItem> _receivedMessages;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Received messages lock
|
/// Received messages lock
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -58,19 +56,19 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <summary>
|
/// <summary>
|
||||||
/// Handlers for when an error happens on the socket
|
/// Handlers for when an error happens on the socket
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected readonly List<Action<Exception>> errorHandlers = new List<Action<Exception>>();
|
protected readonly List<Action<Exception>> errorHandlers = new();
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Handlers for when the socket connection is opened
|
/// Handlers for when the socket connection is opened
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected readonly List<Action> openHandlers = new List<Action>();
|
protected readonly List<Action> openHandlers = new();
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Handlers for when the connection is closed
|
/// Handlers for when the connection is closed
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected readonly List<Action> closeHandlers = new List<Action>();
|
protected readonly List<Action> closeHandlers = new();
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Handlers for when a message is received
|
/// Handlers for when a message is received
|
||||||
/// </summary>
|
/// </summary>
|
||||||
protected readonly List<Action<string>> messageHandlers = new List<Action<string>>();
|
protected readonly List<Action<string>> messageHandlers = new();
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public int Id { get; }
|
public int Id { get; }
|
||||||
@@ -78,9 +76,6 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public string? Origin { get; set; }
|
public string? Origin { get; set; }
|
||||||
|
|
||||||
/// <inheritdoc />
|
|
||||||
public bool Reconnecting { get; set; }
|
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The timestamp this socket has been active for the last time
|
/// The timestamp this socket has been active for the last time
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -97,13 +92,13 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
public Func<string, string>? DataInterpreterString { get; set; }
|
public Func<string, string>? DataInterpreterString { get; set; }
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public string Url { get; }
|
public Uri Uri { get; }
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public bool IsClosed => _socket.State == WebSocketState.Closed;
|
public bool IsClosed => _socket.State == WebSocketState.Closed;
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public bool IsOpen => _socket.State == WebSocketState.Open && !_closing;
|
public bool IsOpen => _socket.State == WebSocketState.Open && !_ctsSource.IsCancellationRequested;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Ssl protocols supported. NOT USED BY THIS IMPLEMENTATION
|
/// Ssl protocols supported. NOT USED BY THIS IMPLEMENTATION
|
||||||
@@ -130,6 +125,9 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public TimeSpan Timeout { get; set; }
|
public TimeSpan Timeout { get; set; }
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public TimeSpan KeepAliveInterval { get; set; }
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public double IncomingKbps
|
public double IncomingKbps
|
||||||
{
|
{
|
||||||
@@ -179,8 +177,8 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// ctor
|
/// ctor
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="log">The log object to use</param>
|
/// <param name="log">The log object to use</param>
|
||||||
/// <param name="url">The url the socket should connect to</param>
|
/// <param name="uri">The uri the socket should connect to</param>
|
||||||
public CryptoExchangeWebSocketClient(Log log, string url) : this(log, url, new Dictionary<string, string>(), new Dictionary<string, string>())
|
public CryptoExchangeWebSocketClient(Log log, Uri uri) : this(log, uri, new Dictionary<string, string>(), new Dictionary<string, string>())
|
||||||
{
|
{
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -188,14 +186,14 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// ctor
|
/// ctor
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="log">The log object to use</param>
|
/// <param name="log">The log object to use</param>
|
||||||
/// <param name="url">The url the socket should connect to</param>
|
/// <param name="uri">The uri the socket should connect to</param>
|
||||||
/// <param name="cookies">Cookies to sent in the socket connection request</param>
|
/// <param name="cookies">Cookies to sent in the socket connection request</param>
|
||||||
/// <param name="headers">Headers to sent in the socket connection request</param>
|
/// <param name="headers">Headers to sent in the socket connection request</param>
|
||||||
public CryptoExchangeWebSocketClient(Log log, string url, IDictionary<string, string> cookies, IDictionary<string, string> headers)
|
public CryptoExchangeWebSocketClient(Log log, Uri uri, IDictionary<string, string> cookies, IDictionary<string, string> headers)
|
||||||
{
|
{
|
||||||
Id = NextStreamId();
|
Id = NextStreamId();
|
||||||
this.log = log;
|
this.log = log;
|
||||||
Url = url;
|
Uri = uri;
|
||||||
this.cookies = cookies;
|
this.cookies = cookies;
|
||||||
this.headers = headers;
|
this.headers = headers;
|
||||||
|
|
||||||
@@ -212,13 +210,18 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public virtual void SetProxy(ApiProxy proxy)
|
public virtual void SetProxy(ApiProxy proxy)
|
||||||
{
|
{
|
||||||
Uri.TryCreate($"{proxy.Host}:{proxy.Port}", UriKind.Absolute, out var uri);
|
_proxy = proxy;
|
||||||
|
|
||||||
|
if (!Uri.TryCreate($"{proxy.Host}:{proxy.Port}", UriKind.Absolute, out var uri))
|
||||||
|
throw new ArgumentException("Proxy settings invalid, {proxy.Host}:{proxy.Port} not a valid URI", nameof(proxy));
|
||||||
|
|
||||||
_socket.Options.Proxy = uri?.Scheme == null
|
_socket.Options.Proxy = uri?.Scheme == null
|
||||||
? _socket.Options.Proxy = new WebProxy(proxy.Host, proxy.Port)
|
? _socket.Options.Proxy = new WebProxy(proxy.Host, proxy.Port)
|
||||||
: _socket.Options.Proxy = new WebProxy
|
: _socket.Options.Proxy = new WebProxy
|
||||||
{
|
{
|
||||||
Address = uri
|
Address = uri
|
||||||
};
|
};
|
||||||
|
|
||||||
if (proxy.Login != null)
|
if (proxy.Login != null)
|
||||||
_socket.Options.Proxy.Credentials = new NetworkCredential(proxy.Login, proxy.Password);
|
_socket.Options.Proxy.Credentials = new NetworkCredential(proxy.Login, proxy.Password);
|
||||||
}
|
}
|
||||||
@@ -229,8 +232,8 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
log.Write(LogLevel.Debug, $"Socket {Id} connecting");
|
log.Write(LogLevel.Debug, $"Socket {Id} connecting");
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
using CancellationTokenSource tcs = new CancellationTokenSource(TimeSpan.FromSeconds(10));
|
using CancellationTokenSource tcs = new(TimeSpan.FromSeconds(10));
|
||||||
await _socket.ConnectAsync(new Uri(Url), tcs.Token).ConfigureAwait(false);
|
await _socket.ConnectAsync(Uri, tcs.Token).ConfigureAwait(false);
|
||||||
|
|
||||||
Handle(openHandlers);
|
Handle(openHandlers);
|
||||||
}
|
}
|
||||||
@@ -239,35 +242,27 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
log.Write(LogLevel.Debug, $"Socket {Id} connection failed: " + e.ToLogString());
|
log.Write(LogLevel.Debug, $"Socket {Id} connection failed: " + e.ToLogString());
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Write(LogLevel.Trace, $"Socket {Id} connection succeeded, starting communication");
|
log.Write(LogLevel.Debug, $"Socket {Id} connected to {Uri}");
|
||||||
_sendTask = Task.Factory.StartNew(SendLoopAsync, default, TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, TaskScheduler.Default).Unwrap();
|
|
||||||
_receiveTask = Task.Factory.StartNew(ReceiveLoopAsync, default, TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, TaskScheduler.Default).Unwrap();
|
|
||||||
if (Timeout != default)
|
|
||||||
_timeoutTask = Task.Run(CheckTimeoutAsync);
|
|
||||||
|
|
||||||
var sw = Stopwatch.StartNew();
|
|
||||||
while (!_startedSent || !_startedReceive)
|
|
||||||
{
|
|
||||||
// Wait for the tasks to have actually started
|
|
||||||
await Task.Delay(10).ConfigureAwait(false);
|
|
||||||
|
|
||||||
if(sw.ElapsedMilliseconds > 5000)
|
|
||||||
{
|
|
||||||
_ = _socket.CloseOutputAsync(WebSocketCloseStatus.NormalClosure, "", default);
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} startup interupted");
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} connected to {Url}");
|
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public virtual async Task ProcessAsync()
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {Id} ProcessAsync started");
|
||||||
|
var sendTask = SendLoopAsync();
|
||||||
|
var receiveTask = ReceiveLoopAsync();
|
||||||
|
var timeoutTask = Timeout != default ? CheckTimeoutAsync() : Task.CompletedTask;
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {Id} processing startup completed");
|
||||||
|
await Task.WhenAll(sendTask, receiveTask, timeoutTask).ConfigureAwait(false);
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {Id} ProcessAsync finished");
|
||||||
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public virtual void Send(string data)
|
public virtual void Send(string data)
|
||||||
{
|
{
|
||||||
if (_closing)
|
if (_ctsSource.IsCancellationRequested)
|
||||||
throw new InvalidOperationException($"Socket {Id} Can't send data when socket is not connected");
|
throw new InvalidOperationException($"Socket {Id} Can't send data when socket is not connected");
|
||||||
|
|
||||||
var bytes = _encoding.GetBytes(data);
|
var bytes = _encoding.GetBytes(data);
|
||||||
@@ -280,36 +275,40 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
public virtual async Task CloseAsync()
|
public virtual async Task CloseAsync()
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} closing");
|
log.Write(LogLevel.Debug, $"Socket {Id} closing");
|
||||||
await CloseInternalAsync(true, true).ConfigureAwait(false);
|
await CloseInternalAsync().ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Internal close method, will wait for each task to complete to gracefully close
|
/// Internal close method
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="waitSend"></param>
|
|
||||||
/// <param name="waitReceive"></param>
|
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
private async Task CloseInternalAsync(bool waitSend, bool waitReceive)
|
private async Task CloseInternalAsync()
|
||||||
{
|
{
|
||||||
if (_closing)
|
if (_closed || _disposed)
|
||||||
return;
|
return;
|
||||||
|
|
||||||
_closing = true;
|
_closed = true;
|
||||||
var tasksToAwait = new List<Task>();
|
|
||||||
if (_socket.State == WebSocketState.Open)
|
|
||||||
tasksToAwait.Add(_socket.CloseOutputAsync(WebSocketCloseStatus.NormalClosure, "Closing", default));
|
|
||||||
|
|
||||||
_ctsSource.Cancel();
|
_ctsSource.Cancel();
|
||||||
_sendEvent.Set();
|
_sendEvent.Set();
|
||||||
if (waitSend)
|
|
||||||
tasksToAwait.Add(_sendTask!);
|
|
||||||
if (waitReceive)
|
|
||||||
tasksToAwait.Add(_receiveTask!);
|
|
||||||
if (_timeoutTask != null)
|
|
||||||
tasksToAwait.Add(_timeoutTask);
|
|
||||||
|
|
||||||
log.Write(LogLevel.Trace, $"Socket {Id} waiting for communication loops to finish");
|
if (_socket.State == WebSocketState.Open)
|
||||||
await Task.WhenAll(tasksToAwait).ConfigureAwait(false);
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
await _socket.CloseOutputAsync(WebSocketCloseStatus.NormalClosure, "Closing", default).ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
catch(Exception)
|
||||||
|
{ } // Can sometimes throw an exception when socket is in aborted state due to timing
|
||||||
|
}
|
||||||
|
else if(_socket.State == WebSocketState.CloseReceived)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
await _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", default).ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
catch (Exception)
|
||||||
|
{ } // Can sometimes throw an exception when socket is in aborted state due to timing
|
||||||
|
}
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} closed");
|
log.Write(LogLevel.Debug, $"Socket {Id} closed");
|
||||||
Handle(closeHandlers);
|
Handle(closeHandlers);
|
||||||
}
|
}
|
||||||
@@ -319,7 +318,11 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public void Dispose()
|
public void Dispose()
|
||||||
{
|
{
|
||||||
|
if (_disposed)
|
||||||
|
return;
|
||||||
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} disposing");
|
log.Write(LogLevel.Debug, $"Socket {Id} disposing");
|
||||||
|
_disposed = true;
|
||||||
_socket.Dispose();
|
_socket.Dispose();
|
||||||
_ctsSource.Dispose();
|
_ctsSource.Dispose();
|
||||||
|
|
||||||
@@ -335,11 +338,13 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} resetting");
|
log.Write(LogLevel.Debug, $"Socket {Id} resetting");
|
||||||
_ctsSource = new CancellationTokenSource();
|
_ctsSource = new CancellationTokenSource();
|
||||||
_closing = false;
|
|
||||||
|
|
||||||
while (_sendBuffer.TryDequeue(out _)) { } // Clear send buffer
|
while (_sendBuffer.TryDequeue(out _)) { } // Clear send buffer
|
||||||
|
|
||||||
_socket = CreateSocket();
|
_socket = CreateSocket();
|
||||||
|
if (_proxy != null)
|
||||||
|
SetProxy(_proxy);
|
||||||
|
_closed = false;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -355,7 +360,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
socket.Options.Cookies = cookieContainer;
|
socket.Options.Cookies = cookieContainer;
|
||||||
foreach (var header in headers)
|
foreach (var header in headers)
|
||||||
socket.Options.SetRequestHeader(header.Key, header.Value);
|
socket.Options.SetRequestHeader(header.Key, header.Value);
|
||||||
socket.Options.KeepAliveInterval = TimeSpan.FromSeconds(10);
|
socket.Options.KeepAliveInterval = KeepAliveInterval;
|
||||||
socket.Options.SetBuffer(65536, 65536); // Setting it to anything bigger than 65536 throws an exception in .net framework
|
socket.Options.SetBuffer(65536, 65536); // Setting it to anything bigger than 65536 throws an exception in .net framework
|
||||||
return socket;
|
return socket;
|
||||||
}
|
}
|
||||||
@@ -366,17 +371,16 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
private async Task SendLoopAsync()
|
private async Task SendLoopAsync()
|
||||||
{
|
{
|
||||||
_startedSent = true;
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
while (true)
|
while (true)
|
||||||
{
|
{
|
||||||
if (_closing)
|
if (_ctsSource.IsCancellationRequested)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
await _sendEvent.WaitAsync().ConfigureAwait(false);
|
await _sendEvent.WaitAsync().ConfigureAwait(false);
|
||||||
|
|
||||||
if (_closing)
|
if (_ctsSource.IsCancellationRequested)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
while (_sendBuffer.TryDequeue(out var data))
|
while (_sendBuffer.TryDequeue(out var data))
|
||||||
@@ -388,7 +392,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
while (MessagesSentLastSecond() >= RatelimitPerSecond)
|
while (MessagesSentLastSecond() >= RatelimitPerSecond)
|
||||||
{
|
{
|
||||||
start ??= DateTime.UtcNow;
|
start ??= DateTime.UtcNow;
|
||||||
await Task.Delay(10).ConfigureAwait(false);
|
await Task.Delay(50).ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (start != null)
|
if (start != null)
|
||||||
@@ -410,7 +414,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
// Connection closed unexpectedly, .NET framework
|
// Connection closed unexpectedly, .NET framework
|
||||||
Handle(errorHandlers, ioe);
|
Handle(errorHandlers, ioe);
|
||||||
_ = Task.Run(async () => await CloseInternalAsync(false, true).ConfigureAwait(false));
|
await CloseInternalAsync().ConfigureAwait(false);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -427,7 +431,6 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
finally
|
finally
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Trace, $"Socket {Id} Send loop finished");
|
log.Write(LogLevel.Trace, $"Socket {Id} Send loop finished");
|
||||||
_startedSent = false;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -437,14 +440,13 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
private async Task ReceiveLoopAsync()
|
private async Task ReceiveLoopAsync()
|
||||||
{
|
{
|
||||||
_startedReceive = true;
|
|
||||||
var buffer = new ArraySegment<byte>(new byte[65536]);
|
var buffer = new ArraySegment<byte>(new byte[65536]);
|
||||||
var received = 0;
|
var received = 0;
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
while (true)
|
while (true)
|
||||||
{
|
{
|
||||||
if (_closing)
|
if (_ctsSource.IsCancellationRequested)
|
||||||
break;
|
break;
|
||||||
|
|
||||||
MemoryStream? memoryStream = null;
|
MemoryStream? memoryStream = null;
|
||||||
@@ -468,7 +470,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
// Connection closed unexpectedly
|
// Connection closed unexpectedly
|
||||||
Handle(errorHandlers, wse);
|
Handle(errorHandlers, wse);
|
||||||
_ = Task.Run(async () => await CloseInternalAsync(true, true).ConfigureAwait(false));
|
await CloseInternalAsync().ConfigureAwait(false);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -476,7 +478,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
// Connection closed unexpectedly
|
// Connection closed unexpectedly
|
||||||
log.Write(LogLevel.Debug, $"Socket {Id} received `Close` message");
|
log.Write(LogLevel.Debug, $"Socket {Id} received `Close` message");
|
||||||
_ = Task.Run(async () => await CloseInternalAsync(true, true).ConfigureAwait(false));
|
await CloseInternalAsync().ConfigureAwait(false);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -515,7 +517,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (receiveResult == null || _closing)
|
if (receiveResult == null || _ctsSource.IsCancellationRequested)
|
||||||
{
|
{
|
||||||
// Error during receiving or cancellation requested, stop.
|
// Error during receiving or cancellation requested, stop.
|
||||||
break;
|
break;
|
||||||
@@ -547,7 +549,6 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
finally
|
finally
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Trace, $"Socket {Id} Receive loop finished");
|
log.Write(LogLevel.Trace, $"Socket {Id} Receive loop finished");
|
||||||
_startedReceive = false;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -615,7 +616,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
while (true)
|
while (true)
|
||||||
{
|
{
|
||||||
if (_closing)
|
if (_ctsSource.IsCancellationRequested)
|
||||||
return;
|
return;
|
||||||
|
|
||||||
if (DateTime.UtcNow - LastActionTime > Timeout)
|
if (DateTime.UtcNow - LastActionTime > Timeout)
|
||||||
|
|||||||
@@ -0,0 +1,47 @@
|
|||||||
|
using CryptoExchange.Net.Objects;
|
||||||
|
using Newtonsoft.Json.Linq;
|
||||||
|
using System;
|
||||||
|
using System.Threading;
|
||||||
|
|
||||||
|
namespace CryptoExchange.Net.Sockets
|
||||||
|
{
|
||||||
|
internal class PendingRequest
|
||||||
|
{
|
||||||
|
public Func<JToken, bool> Handler { get; }
|
||||||
|
public JToken? Result { get; private set; }
|
||||||
|
public bool Completed { get; private set; }
|
||||||
|
public AsyncResetEvent Event { get; }
|
||||||
|
public TimeSpan Timeout { get; }
|
||||||
|
|
||||||
|
private CancellationTokenSource cts;
|
||||||
|
|
||||||
|
public PendingRequest(Func<JToken, bool> handler, TimeSpan timeout)
|
||||||
|
{
|
||||||
|
Handler = handler;
|
||||||
|
Event = new AsyncResetEvent(false, false);
|
||||||
|
Timeout = timeout;
|
||||||
|
|
||||||
|
cts = new CancellationTokenSource(timeout);
|
||||||
|
cts.Token.Register(Fail, false);
|
||||||
|
}
|
||||||
|
|
||||||
|
public bool CheckData(JToken data)
|
||||||
|
{
|
||||||
|
if (Handler(data))
|
||||||
|
{
|
||||||
|
Result = data;
|
||||||
|
Completed = true;
|
||||||
|
Event.Set();
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void Fail()
|
||||||
|
{
|
||||||
|
Completed = true;
|
||||||
|
Event.Set();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -65,12 +65,22 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <summary>
|
/// <summary>
|
||||||
/// If connection is made
|
/// If connection is made
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool Connected { get; private set; }
|
public bool Connected => _socket.IsOpen;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The underlying websocket
|
/// The unique ID of the socket
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public IWebsocket Socket { get; set; }
|
public int SocketId => _socket.Id;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The current kilobytes per second of data being received, averaged over the last 3 seconds
|
||||||
|
/// </summary>
|
||||||
|
public double IncomingKbps => _socket.IncomingKbps;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The connection uri
|
||||||
|
/// </summary>
|
||||||
|
public Uri Uri => _socket.Uri;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The API client the connection is for
|
/// The API client the connection is for
|
||||||
@@ -113,7 +123,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
if (pausedActivity != value)
|
if (pausedActivity != value)
|
||||||
{
|
{
|
||||||
pausedActivity = value;
|
pausedActivity = value;
|
||||||
log.Write(LogLevel.Information, $"Socket {Socket.Id} Paused activity: " + value);
|
log.Write(LogLevel.Information, $"Socket {SocketId} Paused activity: " + value);
|
||||||
if(pausedActivity) ActivityPaused?.Invoke();
|
if(pausedActivity) ActivityPaused?.Invoke();
|
||||||
else ActivityUnpaused?.Invoke();
|
else ActivityUnpaused?.Invoke();
|
||||||
}
|
}
|
||||||
@@ -122,13 +132,37 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
|
|
||||||
private bool pausedActivity;
|
private bool pausedActivity;
|
||||||
private readonly List<SocketSubscription> subscriptions;
|
private readonly List<SocketSubscription> subscriptions;
|
||||||
private readonly object subscriptionLock = new object();
|
private readonly object subscriptionLock = new();
|
||||||
|
|
||||||
private bool lostTriggered;
|
private bool lostTriggered;
|
||||||
private readonly Log log;
|
private readonly Log log;
|
||||||
private readonly BaseSocketClient socketClient;
|
private readonly BaseSocketClient socketClient;
|
||||||
|
|
||||||
private readonly List<PendingRequest> pendingRequests;
|
private readonly List<PendingRequest> pendingRequests;
|
||||||
|
private Task? _socketProcessTask;
|
||||||
|
private Task? _socketReconnectTask;
|
||||||
|
private readonly AsyncResetEvent _reconnectWaitEvent;
|
||||||
|
|
||||||
|
private SocketStatus _status;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Status of the socket connection
|
||||||
|
/// </summary>
|
||||||
|
public SocketStatus Status
|
||||||
|
{
|
||||||
|
get => _status;
|
||||||
|
private set
|
||||||
|
{
|
||||||
|
var oldStatus = _status;
|
||||||
|
_status = value;
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} status changed from {oldStatus} to {_status}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The underlying websocket
|
||||||
|
/// </summary>
|
||||||
|
private readonly IWebsocket _socket;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// New socket connection
|
/// New socket connection
|
||||||
@@ -145,14 +179,263 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
pendingRequests = new List<PendingRequest>();
|
pendingRequests = new List<PendingRequest>();
|
||||||
|
|
||||||
subscriptions = new List<SocketSubscription>();
|
subscriptions = new List<SocketSubscription>();
|
||||||
Socket = socket;
|
_socket = socket;
|
||||||
|
|
||||||
|
_reconnectWaitEvent = new AsyncResetEvent(false, true);
|
||||||
|
|
||||||
|
_socket.Timeout = client.ClientOptions.SocketNoDataTimeout;
|
||||||
|
_socket.OnMessage += ProcessMessage;
|
||||||
|
_socket.OnOpen += SocketOnOpen;
|
||||||
|
_socket.OnClose += () => _reconnectWaitEvent.Set();
|
||||||
|
|
||||||
Socket.Timeout = client.ClientOptions.SocketNoDataTimeout;
|
|
||||||
Socket.OnMessage += ProcessMessage;
|
|
||||||
Socket.OnClose += SocketOnClose;
|
|
||||||
Socket.OnOpen += SocketOnOpen;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Connect the websocket and start processing
|
||||||
|
/// </summary>
|
||||||
|
/// <returns></returns>
|
||||||
|
public async Task<bool> ConnectAsync()
|
||||||
|
{
|
||||||
|
var connected = await _socket.ConnectAsync().ConfigureAwait(false);
|
||||||
|
if (connected)
|
||||||
|
{
|
||||||
|
Status = SocketStatus.Connected;
|
||||||
|
_socketReconnectTask = ReconnectWatcherAsync();
|
||||||
|
_socketProcessTask = _socket.ProcessAsync();
|
||||||
|
}
|
||||||
|
|
||||||
|
return connected;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Retrieve the underlying socket
|
||||||
|
/// </summary>
|
||||||
|
/// <returns></returns>
|
||||||
|
public IWebsocket GetSocket()
|
||||||
|
{
|
||||||
|
return _socket;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Trigger a reconnect of the socket connection
|
||||||
|
/// </summary>
|
||||||
|
/// <returns></returns>
|
||||||
|
public async Task TriggerReconnectAsync()
|
||||||
|
{
|
||||||
|
await _socket.CloseAsync().ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Close the connection
|
||||||
|
/// </summary>
|
||||||
|
/// <returns></returns>
|
||||||
|
public async Task CloseAsync()
|
||||||
|
{
|
||||||
|
if (Status == SocketStatus.Closed || Status == SocketStatus.Disposed)
|
||||||
|
return;
|
||||||
|
|
||||||
|
ShouldReconnect = false;
|
||||||
|
if (socketClient.socketConnections.ContainsKey(SocketId))
|
||||||
|
socketClient.socketConnections.TryRemove(SocketId, out _);
|
||||||
|
|
||||||
|
lock (subscriptionLock)
|
||||||
|
{
|
||||||
|
foreach (var subscription in subscriptions)
|
||||||
|
{
|
||||||
|
if (subscription.CancellationTokenRegistration.HasValue)
|
||||||
|
subscription.CancellationTokenRegistration.Value.Dispose();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
while (Status == SocketStatus.Reconnecting)
|
||||||
|
// Wait for reconnecting to finish
|
||||||
|
await Task.Delay(100).ConfigureAwait(false);
|
||||||
|
|
||||||
|
await _socket.CloseAsync().ConfigureAwait(false);
|
||||||
|
if(_socketProcessTask != null)
|
||||||
|
await _socketProcessTask.ConfigureAwait(false);
|
||||||
|
_socket.Dispose();
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Close a subscription on this connection. If all subscriptions on this connection are closed the connection gets closed as well
|
||||||
|
/// </summary>
|
||||||
|
/// <param name="subscription">Subscription to close</param>
|
||||||
|
/// <returns></returns>
|
||||||
|
public async Task CloseAsync(SocketSubscription subscription)
|
||||||
|
{
|
||||||
|
if (Status == SocketStatus.Closing || Status == SocketStatus.Closed || Status == SocketStatus.Disposed)
|
||||||
|
return;
|
||||||
|
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} closing subscription {subscription.Id}");
|
||||||
|
if (subscription.CancellationTokenRegistration.HasValue)
|
||||||
|
subscription.CancellationTokenRegistration.Value.Dispose();
|
||||||
|
|
||||||
|
if (subscription.Confirmed && _socket.IsOpen)
|
||||||
|
await socketClient.UnsubscribeAsync(this, subscription).ConfigureAwait(false);
|
||||||
|
|
||||||
|
bool shouldCloseConnection;
|
||||||
|
lock (subscriptionLock)
|
||||||
|
{
|
||||||
|
if (Status == SocketStatus.Closing)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} already closing");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
shouldCloseConnection = subscriptions.All(r => !r.UserSubscription || r == subscription);
|
||||||
|
if (shouldCloseConnection)
|
||||||
|
Status = SocketStatus.Closing;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (shouldCloseConnection)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} closing as there are no more subscriptions");
|
||||||
|
await CloseAsync().ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
lock (subscriptionLock)
|
||||||
|
subscriptions.Remove(subscription);
|
||||||
|
}
|
||||||
|
|
||||||
|
private async Task ReconnectAsync()
|
||||||
|
{
|
||||||
|
// Fail all pending requests
|
||||||
|
lock (pendingRequests)
|
||||||
|
{
|
||||||
|
foreach (var pendingRequest in pendingRequests.ToList())
|
||||||
|
{
|
||||||
|
pendingRequest.Fail();
|
||||||
|
pendingRequests.Remove(pendingRequest);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (socketClient.ClientOptions.AutoReconnect && ShouldReconnect)
|
||||||
|
{
|
||||||
|
// Should reconnect
|
||||||
|
DisconnectTime = DateTime.UtcNow;
|
||||||
|
log.Write(LogLevel.Warning, $"Socket {SocketId} Connection lost, will try to reconnect");
|
||||||
|
if (!lostTriggered)
|
||||||
|
{
|
||||||
|
lostTriggered = true;
|
||||||
|
_ = Task.Run(() => ConnectionLost?.Invoke());
|
||||||
|
}
|
||||||
|
|
||||||
|
while (ShouldReconnect)
|
||||||
|
{
|
||||||
|
if (ReconnectTry > 0)
|
||||||
|
{
|
||||||
|
// Wait a bit before attempting reconnect
|
||||||
|
await Task.Delay(socketClient.ClientOptions.ReconnectInterval).ConfigureAwait(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!ShouldReconnect)
|
||||||
|
{
|
||||||
|
// Should reconnect changed to false while waiting to reconnect
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
_socket.Reset();
|
||||||
|
if (!await _socket.ConnectAsync().ConfigureAwait(false))
|
||||||
|
{
|
||||||
|
// Reconnect failed
|
||||||
|
ReconnectTry++;
|
||||||
|
ResubscribeTry = 0;
|
||||||
|
if (socketClient.ClientOptions.MaxReconnectTries != null
|
||||||
|
&& ReconnectTry >= socketClient.ClientOptions.MaxReconnectTries)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Warning, $"Socket {SocketId} failed to reconnect after {ReconnectTry} tries, closing");
|
||||||
|
ShouldReconnect = false;
|
||||||
|
|
||||||
|
if (socketClient.socketConnections.ContainsKey(SocketId))
|
||||||
|
socketClient.socketConnections.TryRemove(SocketId, out _);
|
||||||
|
|
||||||
|
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
||||||
|
// Reached max tries, break loop and leave connection closed
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Continue to try again
|
||||||
|
log.Write(LogLevel.Debug, $"Socket {SocketId} failed to reconnect{(socketClient.ClientOptions.MaxReconnectTries != null ? $", try {ReconnectTry}/{socketClient.ClientOptions.MaxReconnectTries}" : "")}, will try again in {socketClient.ClientOptions.ReconnectInterval}");
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Successfully reconnected, start processing
|
||||||
|
Status = SocketStatus.Connected;
|
||||||
|
_socketProcessTask = _socket.ProcessAsync();
|
||||||
|
|
||||||
|
ReconnectTry = 0;
|
||||||
|
var time = DisconnectTime;
|
||||||
|
DisconnectTime = null;
|
||||||
|
|
||||||
|
log.Write(LogLevel.Information, $"Socket {SocketId} reconnected after {DateTime.UtcNow - time}");
|
||||||
|
|
||||||
|
var reconnectResult = await ProcessReconnectAsync().ConfigureAwait(false);
|
||||||
|
if (!reconnectResult)
|
||||||
|
{
|
||||||
|
// Failed to resubscribe everything
|
||||||
|
ResubscribeTry++;
|
||||||
|
DisconnectTime = time;
|
||||||
|
|
||||||
|
if (socketClient.ClientOptions.MaxResubscribeTries != null &&
|
||||||
|
ResubscribeTry >= socketClient.ClientOptions.MaxResubscribeTries)
|
||||||
|
{
|
||||||
|
log.Write(LogLevel.Warning, $"Socket {SocketId} failed to resubscribe after {ResubscribeTry} tries, closing. Last resubscription error: {reconnectResult.Error}");
|
||||||
|
ShouldReconnect = false;
|
||||||
|
|
||||||
|
if (socketClient.socketConnections.ContainsKey(SocketId))
|
||||||
|
socketClient.socketConnections.TryRemove(SocketId, out _);
|
||||||
|
|
||||||
|
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
||||||
|
}
|
||||||
|
else
|
||||||
|
log.Write(LogLevel.Debug, $"Socket {SocketId} resubscribing all subscriptions failed on reconnected socket{(socketClient.ClientOptions.MaxResubscribeTries != null ? $", try {ResubscribeTry}/{socketClient.ClientOptions.MaxResubscribeTries}" : "")}. Error: {reconnectResult.Error}. Disconnecting and reconnecting.");
|
||||||
|
|
||||||
|
// Failed resubscribe, close socket if it is still open
|
||||||
|
if (_socket.IsOpen)
|
||||||
|
await _socket.CloseAsync().ConfigureAwait(false);
|
||||||
|
else
|
||||||
|
DisconnectTime = DateTime.UtcNow;
|
||||||
|
|
||||||
|
// Break out of the loop, the new processing task should reconnect again
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
// Succesfully reconnected
|
||||||
|
log.Write(LogLevel.Information, $"Socket {SocketId} data connection restored.");
|
||||||
|
ResubscribeTry = 0;
|
||||||
|
if (lostTriggered)
|
||||||
|
{
|
||||||
|
lostTriggered = false;
|
||||||
|
_ = Task.Run(() => ConnectionRestored?.Invoke(time.HasValue ? DateTime.UtcNow - time.Value : TimeSpan.FromSeconds(0)));
|
||||||
|
}
|
||||||
|
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
if (!socketClient.ClientOptions.AutoReconnect && ShouldReconnect)
|
||||||
|
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
||||||
|
|
||||||
|
// No reconnecting needed
|
||||||
|
log.Write(LogLevel.Information, $"Socket {SocketId} closed");
|
||||||
|
if (socketClient.socketConnections.ContainsKey(SocketId))
|
||||||
|
socketClient.socketConnections.TryRemove(SocketId, out _);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Dispose the connection
|
||||||
|
/// </summary>
|
||||||
|
public void Dispose()
|
||||||
|
{
|
||||||
|
Status = SocketStatus.Disposed;
|
||||||
|
_socket.Dispose();
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Process a message received by the socket
|
/// Process a message received by the socket
|
||||||
/// </summary>
|
/// </summary>
|
||||||
@@ -160,7 +443,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
private void ProcessMessage(string data)
|
private void ProcessMessage(string data)
|
||||||
{
|
{
|
||||||
var timestamp = DateTime.UtcNow;
|
var timestamp = DateTime.UtcNow;
|
||||||
log.Write(LogLevel.Trace, $"Socket {Socket.Id} received data: " + data);
|
log.Write(LogLevel.Trace, $"Socket {SocketId} received data: " + data);
|
||||||
if (string.IsNullOrEmpty(data)) return;
|
if (string.IsNullOrEmpty(data)) return;
|
||||||
|
|
||||||
var tokenData = data.ToJToken(log);
|
var tokenData = data.ToJToken(log);
|
||||||
@@ -202,22 +485,37 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
|
|
||||||
// Message was not a request response, check data handlers
|
// Message was not a request response, check data handlers
|
||||||
var messageEvent = new MessageEvent(this, tokenData, socketClient.ClientOptions.OutputOriginalData ? data: null, timestamp);
|
var messageEvent = new MessageEvent(this, tokenData, socketClient.ClientOptions.OutputOriginalData ? data: null, timestamp);
|
||||||
if (!HandleData(messageEvent) && !handledResponse)
|
var (handled, userProcessTime) = HandleData(messageEvent);
|
||||||
|
if (!handled && !handledResponse)
|
||||||
{
|
{
|
||||||
if (!socketClient.UnhandledMessageExpected)
|
if (!socketClient.UnhandledMessageExpected)
|
||||||
log.Write(LogLevel.Warning, $"Socket {Socket.Id} Message not handled: " + tokenData);
|
log.Write(LogLevel.Warning, $"Socket {SocketId} Message not handled: " + tokenData);
|
||||||
UnhandledMessage?.Invoke(tokenData);
|
UnhandledMessage?.Invoke(tokenData);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
var total = DateTime.UtcNow - timestamp;
|
||||||
|
if (userProcessTime.TotalMilliseconds > 500)
|
||||||
|
log.Write(LogLevel.Debug, $"Socket {SocketId} message processing slow ({(int)total.TotalMilliseconds}ms, {(int)userProcessTime.TotalMilliseconds}ms user code), consider offloading data handling to another thread. " +
|
||||||
|
"Data from this socket may arrive late or not at all if message processing is continuously slow.");
|
||||||
|
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} message processed in {(int)total.TotalMilliseconds}ms, ({(int)userProcessTime.TotalMilliseconds}ms user code)");
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Add a subscription to this connection
|
/// Add a subscription to this connection
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="subscription"></param>
|
/// <param name="subscription"></param>
|
||||||
public void AddSubscription(SocketSubscription subscription)
|
public bool AddSubscription(SocketSubscription subscription)
|
||||||
{
|
{
|
||||||
lock(subscriptionLock)
|
lock (subscriptionLock)
|
||||||
|
{
|
||||||
|
if (Status != SocketStatus.None && Status != SocketStatus.Connected)
|
||||||
|
return false;
|
||||||
|
|
||||||
subscriptions.Add(subscription);
|
subscriptions.Add(subscription);
|
||||||
|
log.Write(LogLevel.Trace, $"Socket {SocketId} adding new subscription with id {subscription.Id}, total subscriptions on connection: {subscriptions.Count}");
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -246,13 +544,13 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="messageEvent"></param>
|
/// <param name="messageEvent"></param>
|
||||||
/// <returns>True if the data was successfully handled</returns>
|
/// <returns>True if the data was successfully handled</returns>
|
||||||
private bool HandleData(MessageEvent messageEvent)
|
private (bool, TimeSpan) HandleData(MessageEvent messageEvent)
|
||||||
{
|
{
|
||||||
SocketSubscription? currentSubscription = null;
|
SocketSubscription? currentSubscription = null;
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
var handled = false;
|
var handled = false;
|
||||||
var sw = Stopwatch.StartNew();
|
TimeSpan userCodeDuration = TimeSpan.Zero;
|
||||||
|
|
||||||
// Loop the subscriptions to check if any of them signal us that the message is for them
|
// Loop the subscriptions to check if any of them signal us that the message is for them
|
||||||
List<SocketSubscription> subscriptionsCopy;
|
List<SocketSubscription> subscriptionsCopy;
|
||||||
@@ -267,7 +565,10 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
if (socketClient.MessageMatchesHandler(this, messageEvent.JsonData, subscription.Identifier!))
|
if (socketClient.MessageMatchesHandler(this, messageEvent.JsonData, subscription.Identifier!))
|
||||||
{
|
{
|
||||||
handled = true;
|
handled = true;
|
||||||
|
var userSw = Stopwatch.StartNew();
|
||||||
subscription.MessageHandler(messageEvent);
|
subscription.MessageHandler(messageEvent);
|
||||||
|
userSw.Stop();
|
||||||
|
userCodeDuration = userSw.Elapsed;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
@@ -276,24 +577,21 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
handled = true;
|
handled = true;
|
||||||
messageEvent.JsonData = socketClient.ProcessTokenData(messageEvent.JsonData);
|
messageEvent.JsonData = socketClient.ProcessTokenData(messageEvent.JsonData);
|
||||||
|
var userSw = Stopwatch.StartNew();
|
||||||
subscription.MessageHandler(messageEvent);
|
subscription.MessageHandler(messageEvent);
|
||||||
|
userSw.Stop();
|
||||||
|
userCodeDuration = userSw.Elapsed;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
sw.Stop();
|
return (handled, userCodeDuration);
|
||||||
if (sw.ElapsedMilliseconds > 500)
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Socket.Id} message processing slow ({sw.ElapsedMilliseconds}ms), consider offloading data handling to another thread. " +
|
|
||||||
"Data from this socket may arrive late or not at all if message processing is continuously slow.");
|
|
||||||
else
|
|
||||||
log.Write(LogLevel.Trace, $"Socket {Socket.Id} message processed in {sw.ElapsedMilliseconds}ms");
|
|
||||||
return handled;
|
|
||||||
}
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Error, $"Socket {Socket.Id} Exception during message processing\r\nException: {ex.ToLogString()}\r\nData: {messageEvent.JsonData}");
|
log.Write(LogLevel.Error, $"Socket {SocketId} Exception during message processing\r\nException: {ex.ToLogString()}\r\nData: {messageEvent.JsonData}");
|
||||||
currentSubscription?.InvokeExceptionHandler(ex);
|
currentSubscription?.InvokeExceptionHandler(ex);
|
||||||
return false;
|
return (false, TimeSpan.Zero);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -312,7 +610,10 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
pendingRequests.Add(pending);
|
pendingRequests.Add(pending);
|
||||||
}
|
}
|
||||||
Send(obj);
|
var sendOk = Send(obj);
|
||||||
|
if(!sendOk)
|
||||||
|
pending.Fail();
|
||||||
|
|
||||||
return pending.Event.WaitAsync(timeout);
|
return pending.Event.WaitAsync(timeout);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -322,22 +623,30 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <typeparam name="T">The type of the object to send</typeparam>
|
/// <typeparam name="T">The type of the object to send</typeparam>
|
||||||
/// <param name="obj">The object to send</param>
|
/// <param name="obj">The object to send</param>
|
||||||
/// <param name="nullValueHandling">How null values should be serialized</param>
|
/// <param name="nullValueHandling">How null values should be serialized</param>
|
||||||
public virtual void Send<T>(T obj, NullValueHandling nullValueHandling = NullValueHandling.Ignore)
|
public virtual bool Send<T>(T obj, NullValueHandling nullValueHandling = NullValueHandling.Ignore)
|
||||||
{
|
{
|
||||||
if(obj is string str)
|
if(obj is string str)
|
||||||
Send(str);
|
return Send(str);
|
||||||
else
|
else
|
||||||
Send(JsonConvert.SerializeObject(obj, Formatting.None, new JsonSerializerSettings { NullValueHandling = nullValueHandling }));
|
return Send(JsonConvert.SerializeObject(obj, Formatting.None, new JsonSerializerSettings { NullValueHandling = nullValueHandling }));
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Send string data over the websocket connection
|
/// Send string data over the websocket connection
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="data">The data to send</param>
|
/// <param name="data">The data to send</param>
|
||||||
public virtual void Send(string data)
|
public virtual bool Send(string data)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Trace, $"Socket {Socket.Id} sending data: {data}");
|
log.Write(LogLevel.Trace, $"Socket {SocketId} sending data: {data}");
|
||||||
Socket.Send(data);
|
try
|
||||||
|
{
|
||||||
|
_socket.Send(data);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
catch(Exception)
|
||||||
|
{
|
||||||
|
return false;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -347,140 +656,27 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
{
|
{
|
||||||
ReconnectTry = 0;
|
ReconnectTry = 0;
|
||||||
PausedActivity = false;
|
PausedActivity = false;
|
||||||
Connected = true;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
private async Task ReconnectWatcherAsync()
|
||||||
/// Handler for a socket closing. Reconnects the socket if needed, or removes it from the active socket list if not
|
|
||||||
/// </summary>
|
|
||||||
protected virtual void SocketOnClose()
|
|
||||||
{
|
{
|
||||||
lock (pendingRequests)
|
while (true)
|
||||||
{
|
{
|
||||||
foreach(var pendingRequest in pendingRequests.ToList())
|
await _reconnectWaitEvent.WaitAsync().ConfigureAwait(false);
|
||||||
{
|
if (!ShouldReconnect)
|
||||||
pendingRequest.Fail();
|
return;
|
||||||
pendingRequests.Remove(pendingRequest);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (socketClient.ClientOptions.AutoReconnect && ShouldReconnect)
|
Status = SocketStatus.Reconnecting;
|
||||||
{
|
await ReconnectAsync().ConfigureAwait(false);
|
||||||
if (Socket.Reconnecting)
|
|
||||||
return; // Already reconnecting
|
|
||||||
|
|
||||||
Socket.Reconnecting = true;
|
if (!ShouldReconnect)
|
||||||
|
return;
|
||||||
DisconnectTime = DateTime.UtcNow;
|
|
||||||
log.Write(LogLevel.Warning, $"Socket {Socket.Id} Connection lost, will try to reconnect");
|
|
||||||
if (!lostTriggered)
|
|
||||||
{
|
|
||||||
lostTriggered = true;
|
|
||||||
ConnectionLost?.Invoke();
|
|
||||||
}
|
|
||||||
|
|
||||||
Task.Run(async () =>
|
|
||||||
{
|
|
||||||
while (ShouldReconnect)
|
|
||||||
{
|
|
||||||
if (ReconnectTry > 0)
|
|
||||||
{
|
|
||||||
// Wait a bit before attempting reconnect
|
|
||||||
await Task.Delay(socketClient.ClientOptions.ReconnectInterval).ConfigureAwait(false);
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!ShouldReconnect)
|
|
||||||
{
|
|
||||||
// Should reconnect changed to false while waiting to reconnect
|
|
||||||
Socket.Reconnecting = false;
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
Socket.Reset();
|
|
||||||
if (!await Socket.ConnectAsync().ConfigureAwait(false))
|
|
||||||
{
|
|
||||||
ReconnectTry++;
|
|
||||||
ResubscribeTry = 0;
|
|
||||||
if (socketClient.ClientOptions.MaxReconnectTries != null
|
|
||||||
&& ReconnectTry >= socketClient.ClientOptions.MaxReconnectTries)
|
|
||||||
{
|
|
||||||
log.Write(LogLevel.Warning, $"Socket {Socket.Id} failed to reconnect after {ReconnectTry} tries, closing");
|
|
||||||
ShouldReconnect = false;
|
|
||||||
|
|
||||||
if (socketClient.sockets.ContainsKey(Socket.Id))
|
|
||||||
socketClient.sockets.TryRemove(Socket.Id, out _);
|
|
||||||
|
|
||||||
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Socket.Id} failed to reconnect{(socketClient.ClientOptions.MaxReconnectTries != null ? $", try {ReconnectTry}/{socketClient.ClientOptions.MaxReconnectTries}": "")}, will try again in {socketClient.ClientOptions.ReconnectInterval}");
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Successfully reconnected
|
|
||||||
var time = DisconnectTime;
|
|
||||||
DisconnectTime = null;
|
|
||||||
|
|
||||||
log.Write(LogLevel.Information, $"Socket {Socket.Id} reconnected after {DateTime.UtcNow - time}");
|
|
||||||
|
|
||||||
var reconnectResult = await ProcessReconnectAsync().ConfigureAwait(false);
|
|
||||||
if (!reconnectResult)
|
|
||||||
{
|
|
||||||
ResubscribeTry++;
|
|
||||||
DisconnectTime = time;
|
|
||||||
|
|
||||||
if (socketClient.ClientOptions.MaxResubscribeTries != null &&
|
|
||||||
ResubscribeTry >= socketClient.ClientOptions.MaxResubscribeTries)
|
|
||||||
{
|
|
||||||
log.Write(LogLevel.Warning, $"Socket {Socket.Id} failed to resubscribe after {ResubscribeTry} tries, closing. Last resubscription error: {reconnectResult.Error}");
|
|
||||||
ShouldReconnect = false;
|
|
||||||
|
|
||||||
if (socketClient.sockets.ContainsKey(Socket.Id))
|
|
||||||
socketClient.sockets.TryRemove(Socket.Id, out _);
|
|
||||||
|
|
||||||
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
|
||||||
}
|
|
||||||
else
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Socket.Id} resubscribing all subscriptions failed on reconnected socket{(socketClient.ClientOptions.MaxResubscribeTries != null ? $", try {ResubscribeTry}/{socketClient.ClientOptions.MaxResubscribeTries}" : "")}. Error: {reconnectResult.Error}. Disconnecting and reconnecting.");
|
|
||||||
|
|
||||||
if (Socket.IsOpen)
|
|
||||||
await Socket.CloseAsync().ConfigureAwait(false);
|
|
||||||
else
|
|
||||||
DisconnectTime = DateTime.UtcNow;
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
log.Write(LogLevel.Information, $"Socket {Socket.Id} data connection restored.");
|
|
||||||
ResubscribeTry = 0;
|
|
||||||
if (lostTriggered)
|
|
||||||
{
|
|
||||||
lostTriggered = false;
|
|
||||||
_ = Task.Run(() => ConnectionRestored?.Invoke(time.HasValue ? DateTime.UtcNow - time.Value : TimeSpan.FromSeconds(0))).ConfigureAwait(false);
|
|
||||||
}
|
|
||||||
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Socket.Reconnecting = false;
|
|
||||||
});
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
if (!socketClient.ClientOptions.AutoReconnect && ShouldReconnect)
|
|
||||||
_ = Task.Run(() => ConnectionClosed?.Invoke());
|
|
||||||
|
|
||||||
// No reconnecting needed
|
|
||||||
log.Write(LogLevel.Information, $"Socket {Socket.Id} closed");
|
|
||||||
if (socketClient.sockets.ContainsKey(Socket.Id))
|
|
||||||
socketClient.sockets.TryRemove(Socket.Id, out _);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task<CallResult<bool>> ProcessReconnectAsync()
|
private async Task<CallResult<bool>> ProcessReconnectAsync()
|
||||||
{
|
{
|
||||||
if (!Socket.IsOpen)
|
if (!_socket.IsOpen)
|
||||||
return new CallResult<bool>(new WebError("Socket not connected"));
|
return new CallResult<bool>(new WebError("Socket not connected"));
|
||||||
|
|
||||||
if (Authenticated)
|
if (Authenticated)
|
||||||
@@ -489,11 +685,11 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
var authResult = await socketClient.AuthenticateSocketAsync(this).ConfigureAwait(false);
|
var authResult = await socketClient.AuthenticateSocketAsync(this).ConfigureAwait(false);
|
||||||
if (!authResult)
|
if (!authResult)
|
||||||
{
|
{
|
||||||
log.Write(LogLevel.Warning, $"Socket {Socket.Id} authentication failed on reconnected socket. Disconnecting and reconnecting.");
|
log.Write(LogLevel.Warning, $"Socket {SocketId} authentication failed on reconnected socket. Disconnecting and reconnecting.");
|
||||||
return authResult;
|
return authResult;
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Socket.Id} authentication succeeded on reconnected socket.");
|
log.Write(LogLevel.Debug, $"Socket {SocketId} authentication succeeded on reconnected socket.");
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get a list of all subscriptions on the socket
|
// Get a list of all subscriptions on the socket
|
||||||
@@ -504,7 +700,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
// Foreach subscription which is subscribed by a subscription request we will need to resend that request to resubscribe
|
// Foreach subscription which is subscribed by a subscription request we will need to resend that request to resubscribe
|
||||||
for (var i = 0; i < subscriptionList.Count; i += socketClient.ClientOptions.MaxConcurrentResubscriptionsPerSocket)
|
for (var i = 0; i < subscriptionList.Count; i += socketClient.ClientOptions.MaxConcurrentResubscriptionsPerSocket)
|
||||||
{
|
{
|
||||||
if (!Socket.IsOpen)
|
if (!_socket.IsOpen)
|
||||||
return new CallResult<bool>(new WebError("Socket not connected"));
|
return new CallResult<bool>(new WebError("Socket not connected"));
|
||||||
|
|
||||||
var taskList = new List<Task<CallResult<bool>>>();
|
var taskList = new List<Task<CallResult<bool>>>();
|
||||||
@@ -516,10 +712,10 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
return taskList.First(t => !t.Result.Success).Result;
|
return taskList.First(t => !t.Result.Success).Result;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!Socket.IsOpen)
|
if (!_socket.IsOpen)
|
||||||
return new CallResult<bool>(new WebError("Socket not connected"));
|
return new CallResult<bool>(new WebError("Socket not connected"));
|
||||||
|
|
||||||
log.Write(LogLevel.Debug, $"Socket {Socket.Id} all subscription successfully resubscribed on reconnected socket.");
|
log.Write(LogLevel.Debug, $"Socket {SocketId} all subscription successfully resubscribed on reconnected socket.");
|
||||||
return new CallResult<bool>(true);
|
return new CallResult<bool>(true);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -530,100 +726,41 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
|
|
||||||
internal async Task<CallResult<bool>> ResubscribeAsync(SocketSubscription socketSubscription)
|
internal async Task<CallResult<bool>> ResubscribeAsync(SocketSubscription socketSubscription)
|
||||||
{
|
{
|
||||||
if (!Socket.IsOpen)
|
if (!_socket.IsOpen)
|
||||||
return new CallResult<bool>(new UnknownError("Socket is not connected"));
|
return new CallResult<bool>(new UnknownError("Socket is not connected"));
|
||||||
|
|
||||||
return await socketClient.SubscribeAndWaitAsync(this, socketSubscription.Request!, socketSubscription).ConfigureAwait(false);
|
return await socketClient.SubscribeAndWaitAsync(this, socketSubscription.Request!, socketSubscription).ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Close the connection
|
/// Status of the socket connection
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <returns></returns>
|
public enum SocketStatus
|
||||||
public async Task CloseAsync()
|
|
||||||
{
|
{
|
||||||
Connected = false;
|
/// <summary>
|
||||||
ShouldReconnect = false;
|
/// None/Initial
|
||||||
if (socketClient.sockets.ContainsKey(Socket.Id))
|
/// </summary>
|
||||||
socketClient.sockets.TryRemove(Socket.Id, out _);
|
None,
|
||||||
|
/// <summary>
|
||||||
lock (subscriptionLock)
|
/// Connected
|
||||||
{
|
/// </summary>
|
||||||
foreach (var subscription in subscriptions)
|
Connected,
|
||||||
{
|
/// <summary>
|
||||||
if (subscription.CancellationTokenRegistration.HasValue)
|
/// Reconnecting
|
||||||
subscription.CancellationTokenRegistration.Value.Dispose();
|
/// </summary>
|
||||||
}
|
Reconnecting,
|
||||||
}
|
/// <summary>
|
||||||
await Socket.CloseAsync().ConfigureAwait(false);
|
/// Closing
|
||||||
Socket.Dispose();
|
/// </summary>
|
||||||
}
|
Closing,
|
||||||
|
/// <summary>
|
||||||
/// <summary>
|
/// Closed
|
||||||
/// Close a subscription on this connection. If all subscriptions on this connection are closed the connection gets closed as well
|
/// </summary>
|
||||||
/// </summary>
|
Closed,
|
||||||
/// <param name="subscription">Subscription to close</param>
|
/// <summary>
|
||||||
/// <returns></returns>
|
/// Disposed
|
||||||
public async Task CloseAsync(SocketSubscription subscription)
|
/// </summary>
|
||||||
{
|
Disposed
|
||||||
if (!Socket.IsOpen)
|
|
||||||
return;
|
|
||||||
|
|
||||||
if (subscription.CancellationTokenRegistration.HasValue)
|
|
||||||
subscription.CancellationTokenRegistration.Value.Dispose();
|
|
||||||
|
|
||||||
if (subscription.Confirmed)
|
|
||||||
await socketClient.UnsubscribeAsync(this, subscription).ConfigureAwait(false);
|
|
||||||
|
|
||||||
bool shouldCloseConnection;
|
|
||||||
lock (subscriptionLock)
|
|
||||||
shouldCloseConnection = !subscriptions.Any(r => r.UserSubscription && subscription != r);
|
|
||||||
|
|
||||||
if (shouldCloseConnection)
|
|
||||||
await CloseAsync().ConfigureAwait(false);
|
|
||||||
|
|
||||||
lock (subscriptionLock)
|
|
||||||
subscriptions.Remove(subscription);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
internal class PendingRequest
|
|
||||||
{
|
|
||||||
public Func<JToken, bool> Handler { get; }
|
|
||||||
public JToken? Result { get; private set; }
|
|
||||||
public bool Completed { get; private set; }
|
|
||||||
public AsyncResetEvent Event { get; }
|
|
||||||
public TimeSpan Timeout { get; }
|
|
||||||
|
|
||||||
private CancellationTokenSource cts;
|
|
||||||
|
|
||||||
public PendingRequest(Func<JToken, bool> handler, TimeSpan timeout)
|
|
||||||
{
|
|
||||||
Handler = handler;
|
|
||||||
Event = new AsyncResetEvent(false, false);
|
|
||||||
Timeout = timeout;
|
|
||||||
|
|
||||||
cts = new CancellationTokenSource(timeout);
|
|
||||||
cts.Token.Register(Fail, false);
|
|
||||||
}
|
|
||||||
|
|
||||||
public bool CheckData(JToken data)
|
|
||||||
{
|
|
||||||
if (Handler(data))
|
|
||||||
{
|
|
||||||
Result = data;
|
|
||||||
Completed = true;
|
|
||||||
Event.Set();
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
|
|
||||||
public void Fail()
|
|
||||||
{
|
|
||||||
Completed = true;
|
|
||||||
Event.Set();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -72,7 +72,7 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// <summary>
|
/// <summary>
|
||||||
/// The id of the socket
|
/// The id of the socket
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public int SocketId => connection.Socket.Id;
|
public int SocketId => connection.SocketId;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// The id of the subscription
|
/// The id of the subscription
|
||||||
@@ -103,9 +103,9 @@ namespace CryptoExchange.Net.Sockets
|
|||||||
/// Close the socket to cause a reconnect
|
/// Close the socket to cause a reconnect
|
||||||
/// </summary>
|
/// </summary>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
internal Task ReconnectAsync()
|
public Task ReconnectAsync()
|
||||||
{
|
{
|
||||||
return connection.Socket.CloseAsync();
|
return connection.TriggerReconnectAsync();
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
|
|||||||
@@ -1,24 +1,25 @@
|
|||||||
using System.Collections.Generic;
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
using CryptoExchange.Net.Interfaces;
|
using CryptoExchange.Net.Interfaces;
|
||||||
using CryptoExchange.Net.Logging;
|
using CryptoExchange.Net.Logging;
|
||||||
|
|
||||||
namespace CryptoExchange.Net.Sockets
|
namespace CryptoExchange.Net.Sockets
|
||||||
{
|
{
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Default weboscket factory implementation
|
/// Default websocket factory implementation
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public class WebsocketFactory : IWebsocketFactory
|
public class WebsocketFactory : IWebsocketFactory
|
||||||
{
|
{
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public IWebsocket CreateWebsocket(Log log, string url)
|
public IWebsocket CreateWebsocket(Log log, string url)
|
||||||
{
|
{
|
||||||
return new CryptoExchangeWebSocketClient(log, url);
|
return new CryptoExchangeWebSocketClient(log, new Uri(url));
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public IWebsocket CreateWebsocket(Log log, string url, IDictionary<string, string> cookies, IDictionary<string, string> headers)
|
public IWebsocket CreateWebsocket(Log log, string url, IDictionary<string, string> cookies, IDictionary<string, string> headers)
|
||||||
{
|
{
|
||||||
return new CryptoExchangeWebSocketClient(log, url, cookies, headers);
|
return new CryptoExchangeWebSocketClient(log, new Uri(url), cookies, headers);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,6 +18,39 @@ I develop and maintain this package on my own for free in my spare time. Donatio
|
|||||||
Alternatively, sponsor me on Github using [Github Sponsors](https://github.com/sponsors/JKorf)
|
Alternatively, sponsor me on Github using [Github Sponsors](https://github.com/sponsors/JKorf)
|
||||||
|
|
||||||
## Release notes
|
## Release notes
|
||||||
|
* Version 5.1.12 - 12 Jun 2022
|
||||||
|
* Changed time sync so requests no longer wait for it to complete unless it's the first time
|
||||||
|
* Made log client options changable after client creation
|
||||||
|
* Fixed proxy setting not used when reconnecting socket
|
||||||
|
* Changed MaxSocketConnections to a client options
|
||||||
|
* Updated socket reconnection logic
|
||||||
|
|
||||||
|
* Version 5.1.12 - 12 Jun 2022
|
||||||
|
* Changed time sync so requests no longer wait for it to complete unless it's the first time
|
||||||
|
* Made log client options changable after client creation
|
||||||
|
* Fixed proxy setting not used when reconnecting socket
|
||||||
|
* Updated socket reconnection logic
|
||||||
|
|
||||||
|
* Version 5.1.11 - 24 May 2022
|
||||||
|
* Added KeepAliveInterval setting
|
||||||
|
* Fixed port not being copied when setting parameters on request
|
||||||
|
* Fixed inconsistent PackageReference casing in csproj
|
||||||
|
|
||||||
|
* Version 5.1.10 - 22 May 2022
|
||||||
|
* Fixed order book reconnecting while Diposed
|
||||||
|
* Fixed exception when disposing socket client while reconnecting
|
||||||
|
* Added additional null/default checking in DateTimeConverter
|
||||||
|
* Changed ConnectionLost subscription event to run in seperate task to prevent exception/longer operations from intervering with reconnecting
|
||||||
|
|
||||||
|
* Version 5.1.9 - 08 May 2022
|
||||||
|
* Added latency to the timesync calculation
|
||||||
|
* Small fix for exception in socket close handling
|
||||||
|
|
||||||
|
* Version 5.1.8 - 01 May 2022
|
||||||
|
* Cleanup socket code, fixed an issue which could cause connections to never reconnect when connection was lost
|
||||||
|
* Added support for sending requests which expect an empty response
|
||||||
|
* Fixed issue with the DateTimeConverter date interpretation
|
||||||
|
|
||||||
* Version 5.1.7 - 14 Apr 2022
|
* Version 5.1.7 - 14 Apr 2022
|
||||||
* Moved some Rest parameters from BaseRestClient to RestApiClient to allow different implementations for sub clients
|
* Moved some Rest parameters from BaseRestClient to RestApiClient to allow different implementations for sub clients
|
||||||
|
|
||||||
|
|||||||
+4
-1
@@ -61,4 +61,7 @@ var client = new BinanceClient(new BinanceClientOptions
|
|||||||
BaseAddress = BinanceApiAddresses.TestNet.UsdFuturesRestClientAddress
|
BaseAddress = BinanceApiAddresses.TestNet.UsdFuturesRestClientAddress
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
```
|
```
|
||||||
|
|
||||||
|
### How are timezones handled / Timestamps are off by xx
|
||||||
|
Exchange API's treat all timestamps as UTC, both incoming and outgoing. The client libraries do no conversion so be sure to use UTC as well.
|
||||||
Reference in New Issue
Block a user