1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-13 09:23:04 +00:00

Compare commits

..

4 Commits

Author SHA1 Message Date
JKorf 99c331b389 Updated version 2022-07-10 16:41:14 +02:00
JKorf 70f8bd203a Finished up websocket refactoring 2022-07-10 16:36:00 +02:00
JKorf 89b517c936 wip 2022-07-08 20:24:58 +02:00
JKorf 91e33cc42c wip 2022-07-07 22:17:55 +02:00
16 changed files with 93 additions and 191 deletions
@@ -42,7 +42,7 @@ namespace CryptoExchange.Net.UnitTests
socket.CanConnect = canConnect; socket.CanConnect = canConnect;
//act //act
var connectResult = client.ConnectSocketSub(new SocketConnection(client, null, socket, null)); var connectResult = client.ConnectSocketSub(new SocketConnection(client, null, socket));
//assert //assert
Assert.IsTrue(connectResult.Success == canConnect); Assert.IsTrue(connectResult.Success == canConnect);
@@ -57,10 +57,10 @@ namespace CryptoExchange.Net.UnitTests
socket.ShouldReconnect = true; socket.ShouldReconnect = true;
socket.CanConnect = true; socket.CanConnect = true;
socket.DisconnectTime = DateTime.UtcNow; socket.DisconnectTime = DateTime.UtcNow;
var sub = new SocketConnection(client, null, socket, null); var sub = new SocketConnection(client, null, socket);
var rstEvent = new ManualResetEvent(false); var rstEvent = new ManualResetEvent(false);
JToken result = null; JToken result = null;
sub.AddSubscription(SocketSubscription.CreateForIdentifier(10, "TestHandler", true, false, (messageEvent) => sub.AddSubscription(SocketSubscription.CreateForIdentifier(10, "TestHandler", true, (messageEvent) =>
{ {
result = messageEvent.JsonData; result = messageEvent.JsonData;
rstEvent.Set(); rstEvent.Set();
@@ -85,10 +85,10 @@ namespace CryptoExchange.Net.UnitTests
socket.ShouldReconnect = true; socket.ShouldReconnect = true;
socket.CanConnect = true; socket.CanConnect = true;
socket.DisconnectTime = DateTime.UtcNow; socket.DisconnectTime = DateTime.UtcNow;
var sub = new SocketConnection(client, null, socket, null); var sub = new SocketConnection(client, null, socket);
var rstEvent = new ManualResetEvent(false); var rstEvent = new ManualResetEvent(false);
string original = null; string original = null;
sub.AddSubscription(SocketSubscription.CreateForIdentifier(10, "TestHandler", true, false, (messageEvent) => sub.AddSubscription(SocketSubscription.CreateForIdentifier(10, "TestHandler", true, (messageEvent) =>
{ {
original = messageEvent.OriginalData; original = messageEvent.OriginalData;
rstEvent.Set(); rstEvent.Set();
@@ -103,6 +103,34 @@ namespace CryptoExchange.Net.UnitTests
Assert.IsTrue(original == (enabled ? "{\"property\": 123}" : null)); Assert.IsTrue(original == (enabled ? "{\"property\": 123}" : null));
} }
[TestCase]
public void DisconnectedSocket_Should_Reconnect()
{
// arrange
bool reconnected = false;
var client = new TestSocketClient(new TestOptions() { ReconnectInterval = TimeSpan.Zero, LogLevel = LogLevel.Debug });
var socket = client.CreateSocket();
socket.ShouldReconnect = true;
socket.CanConnect = true;
socket.DisconnectTime = DateTime.UtcNow;
var sub = new SocketConnection(client, null, socket);
sub.ShouldReconnect = true;
client.ConnectSocketSub(sub);
var rstEvent = new ManualResetEvent(false);
sub.ConnectionRestored += (a) =>
{
reconnected = true;
rstEvent.Set();
};
// act
socket.InvokeClose();
rstEvent.WaitOne(1000);
// assert
Assert.IsTrue(reconnected);
}
[TestCase()] [TestCase()]
public void UnsubscribingStream_Should_CloseTheSocket() public void UnsubscribingStream_Should_CloseTheSocket()
{ {
@@ -110,11 +138,9 @@ namespace CryptoExchange.Net.UnitTests
var client = new TestSocketClient(new TestOptions() { ReconnectInterval = TimeSpan.Zero, LogLevel = LogLevel.Debug }); var client = new TestSocketClient(new TestOptions() { ReconnectInterval = TimeSpan.Zero, LogLevel = LogLevel.Debug });
var socket = client.CreateSocket(); var socket = client.CreateSocket();
socket.CanConnect = true; socket.CanConnect = true;
var sub = new SocketConnection(client, null, socket, null); var sub = new SocketConnection(client, null, socket);
client.ConnectSocketSub(sub); client.ConnectSocketSub(sub);
var us = SocketSubscription.CreateForIdentifier(10, "Test", true, false, (e) => { }); var ups = new UpdateSubscription(sub, SocketSubscription.CreateForIdentifier(10, "Test", true, (e) => {}));
var ups = new UpdateSubscription(sub, us);
sub.AddSubscription(us);
// act // act
client.UnsubscribeAsync(ups).Wait(); client.UnsubscribeAsync(ups).Wait();
@@ -132,8 +158,8 @@ namespace CryptoExchange.Net.UnitTests
var socket2 = client.CreateSocket(); var socket2 = client.CreateSocket();
socket1.CanConnect = true; socket1.CanConnect = true;
socket2.CanConnect = true; socket2.CanConnect = true;
var sub1 = new SocketConnection(client, null, socket1, null); var sub1 = new SocketConnection(client, null, socket1);
var sub2 = new SocketConnection(client, null, socket2, null); var sub2 = new SocketConnection(client, null, socket2);
client.ConnectSocketSub(sub1); client.ConnectSocketSub(sub1);
client.ConnectSocketSub(sub2); client.ConnectSocketSub(sub2);
@@ -152,7 +178,7 @@ namespace CryptoExchange.Net.UnitTests
var client = new TestSocketClient(new TestOptions() { ReconnectInterval = TimeSpan.Zero, LogLevel = LogLevel.Debug }); var client = new TestSocketClient(new TestOptions() { ReconnectInterval = TimeSpan.Zero, LogLevel = LogLevel.Debug });
var socket = client.CreateSocket(); var socket = client.CreateSocket();
socket.CanConnect = false; socket.CanConnect = false;
var sub = new SocketConnection(client, null, socket, null); var sub = new SocketConnection(client, null, socket);
// act // act
var connectResult = client.ConnectSocketSub(sub); var connectResult = client.ConnectSocketSub(sub);
@@ -13,15 +13,9 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
public bool Connected { get; set; } public bool Connected { get; set; }
public event Action OnClose; public event Action OnClose;
#pragma warning disable 0067
public event Action OnReconnected;
public event Action OnReconnecting;
#pragma warning restore 0067
public event Action<string> OnMessage; public event Action<string> OnMessage;
public event Action<Exception> OnError; public event Action<Exception> OnError;
public event Action OnOpen; public event Action OnOpen;
public Func<Task<Uri>> GetReconnectionUrl { get; set; }
public int Id { get; } public int Id { get; }
public bool ShouldReconnect { get; set; } public bool ShouldReconnect { get; set; }
@@ -99,7 +93,6 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
{ {
Connected = false; Connected = false;
DisconnectTime = DateTime.UtcNow; DisconnectTime = DateTime.UtcNow;
Reconnecting = true;
OnClose?.Invoke(); OnClose?.Invoke();
} }
@@ -122,6 +115,11 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
{ {
OnError?.Invoke(error); OnError?.Invoke(error);
} }
public Task ReconnectAsync() => Task.CompletedTask;
public async Task ProcessAsync()
{
while (Connected)
await Task.Delay(50);
}
} }
} }
@@ -22,13 +22,13 @@ namespace CryptoExchange.Net.UnitTests.TestImplementations
{ {
SubClient = new TestSubSocketClient(exchangeOptions, exchangeOptions.SubOptions); SubClient = new TestSubSocketClient(exchangeOptions, exchangeOptions.SubOptions);
SocketFactory = new Mock<IWebsocketFactory>().Object; SocketFactory = new Mock<IWebsocketFactory>().Object;
Mock.Get(SocketFactory).Setup(f => f.CreateWebsocket(It.IsAny<Log>(), It.IsAny<WebSocketParameters>())).Returns(new TestSocket()); Mock.Get(SocketFactory).Setup(f => f.CreateWebsocket(It.IsAny<Log>(), It.IsAny<string>())).Returns(new TestSocket());
} }
public TestSocket CreateSocket() public TestSocket CreateSocket()
{ {
Mock.Get(SocketFactory).Setup(f => f.CreateWebsocket(It.IsAny<Log>(), It.IsAny<WebSocketParameters>())).Returns(new TestSocket()); Mock.Get(SocketFactory).Setup(f => f.CreateWebsocket(It.IsAny<Log>(), It.IsAny<string>())).Returns(new TestSocket());
return (TestSocket)CreateSocket("https://localhost:123/"); return (TestSocket)CreateSocket("123");
} }
public CallResult<bool> ConnectSocketSub(SocketConnection sub) public CallResult<bool> ConnectSocketSub(SocketConnection sub)
+1 -1
View File
@@ -320,7 +320,7 @@ namespace CryptoExchange.Net
responseStream.Close(); responseStream.Close();
response.Close(); response.Close();
var parseResult = ValidateJson(data); var parseResult = ValidateJson(data);
var error = parseResult.Success ? ParseErrorResponse(parseResult.Data) : new ServerError(data)!; var error = parseResult.Success ? ParseErrorResponse(parseResult.Data) : parseResult.Error!;
if(error.Code == null || error.Code == 0) if(error.Code == null || error.Code == 0)
error.Code = (int)response.StatusCode; error.Code = (int)response.StatusCode;
return new WebCallResult<T>(statusCode, headers, sw.Elapsed, data, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, error); return new WebCallResult<T>(statusCode, headers, sw.Elapsed, data, request.Uri.ToString(), request.Content, request.Method, request.GetHeaders(), default, error);
+1 -15
View File
@@ -392,14 +392,11 @@ namespace CryptoExchange.Net
if (!authenticated || socket.Authenticated) if (!authenticated || socket.Authenticated)
return new CallResult<bool>(true); return new CallResult<bool>(true);
log.Write(LogLevel.Debug, $"Attempting to authenticate {socket.SocketId}");
var result = await AuthenticateSocketAsync(socket).ConfigureAwait(false); var result = await AuthenticateSocketAsync(socket).ConfigureAwait(false);
if (!result) if (!result)
{ {
await socket.CloseAsync().ConfigureAwait(false);
log.Write(LogLevel.Warning, $"Socket {socket.SocketId} authentication failed"); log.Write(LogLevel.Warning, $"Socket {socket.SocketId} authentication failed");
if(socket.Connected)
await socket.CloseAsync().ConfigureAwait(false);
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);
} }
@@ -543,17 +540,6 @@ namespace CryptoExchange.Net
return Task.FromResult(new CallResult<string?>(address)); return Task.FromResult(new CallResult<string?>(address));
} }
/// <summary>
/// Get the url to reconnect to after losing a connection
/// </summary>
/// <param name="apiClient"></param>
/// <param name="connection"></param>
/// <returns></returns>
public virtual Task<Uri?> GetReconnectUriAsync(SocketApiClient apiClient, SocketConnection connection)
{
return Task.FromResult<Uri?>(connection.ConnectionUri);
}
/// <summary> /// <summary>
/// Gets a connection for a new subscription or query. Can be an existing if there are open position or a new one. /// Gets a connection for a new subscription or query. Can be an existing if there are open position or a new one.
/// </summary> /// </summary>
@@ -132,7 +132,7 @@ namespace CryptoExchange.Net.Converters
public override void WriteJson(JsonWriter writer, object? value, JsonSerializer serializer) public override void WriteJson(JsonWriter writer, object? value, JsonSerializer serializer)
{ {
var stringValue = GetString(value); var stringValue = GetString(value);
writer.WriteValue(stringValue); writer.WriteRawValue(stringValue);
} }
} }
} }
+4 -4
View File
@@ -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.2.4</PackageVersion> <PackageVersion>5.2.0</PackageVersion>
<AssemblyVersion>5.2.4</AssemblyVersion> <AssemblyVersion>5.2.0</AssemblyVersion>
<FileVersion>5.2.4</FileVersion> <FileVersion>5.2.0</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.2.4 - Added handling of PlatformNotSupportedException when trying to use websocket from WebAssembly, Changed DataEvent to have a public constructor for testing purposes, Fixed EnumConverter serializing values without proper quotes, Fixed websocket connection reconnecting too quickly when resubscribing/reauthenticating fails</PackageReleaseNotes> <PackageReleaseNotes>5.2.0 - Refactored websocket code, removed some clutter and simplified, Added ReconnectAsync and GetSubscriptionsState methods on socket clients</PackageReleaseNotes>
<Nullable>enable</Nullable> <Nullable>enable</Nullable>
<LangVersion>9.0</LangVersion> <LangVersion>9.0</LangVersion>
<PackageLicenseExpression>MIT</PackageLicenseExpression> <PackageLicenseExpression>MIT</PackageLicenseExpression>
@@ -1,5 +1,4 @@
using CryptoExchange.Net.Objects; using CryptoExchange.Net.Objects;
using CryptoExchange.Net.Sockets;
using System; using System;
using System.Security.Authentication; using System.Security.Authentication;
using System.Text; using System.Text;
@@ -36,10 +35,6 @@ namespace CryptoExchange.Net.Interfaces
/// Websocket has reconnected to the server /// Websocket has reconnected to the server
/// </summary> /// </summary>
event Action OnReconnected; event Action OnReconnected;
/// <summary>
/// Get reconntion url
/// </summary>
Func<Task<Uri?>> GetReconnectionUrl { get; set; }
/// <summary> /// <summary>
/// Unique id for this socket /// Unique id for this socket
@@ -9,7 +9,6 @@ using System.IO;
using System.Linq; using System.Linq;
using System.Net; using System.Net;
using System.Net.WebSockets; using System.Net.WebSockets;
using System.Runtime.InteropServices;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
@@ -34,6 +33,7 @@ namespace CryptoExchange.Net.Sockets
private readonly AsyncResetEvent _sendEvent; private readonly AsyncResetEvent _sendEvent;
private readonly ConcurrentQueue<byte[]> _sendBuffer; private readonly ConcurrentQueue<byte[]> _sendBuffer;
private readonly SemaphoreSlim _closeSem; private readonly SemaphoreSlim _closeSem;
private readonly WebSocketParameters _parameters;
private readonly List<DateTime> _outgoingMessages; private readonly List<DateTime> _outgoingMessages;
private ClientWebSocket _socket; private ClientWebSocket _socket;
@@ -44,7 +44,6 @@ namespace CryptoExchange.Net.Sockets
private bool _stopRequested; private bool _stopRequested;
private bool _disposed; private bool _disposed;
private ProcessState _processState; private ProcessState _processState;
private DateTime _lastReconnectTime;
/// <summary> /// <summary>
@@ -65,16 +64,13 @@ namespace CryptoExchange.Net.Sockets
/// <inheritdoc /> /// <inheritdoc />
public int Id { get; } public int Id { get; }
/// <inheritdoc />
public WebSocketParameters Parameters { get; }
/// <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>
public DateTime LastActionTime { get; private set; } public DateTime LastActionTime { get; private set; }
/// <inheritdoc /> /// <inheritdoc />
public Uri Uri => Parameters.Uri; public Uri Uri => _parameters.Uri;
/// <inheritdoc /> /// <inheritdoc />
public bool IsClosed => _socket.State == WebSocketState.Closed; public bool IsClosed => _socket.State == WebSocketState.Closed;
@@ -111,8 +107,6 @@ namespace CryptoExchange.Net.Sockets
public event Action? OnReconnecting; public event Action? OnReconnecting;
/// <inheritdoc /> /// <inheritdoc />
public event Action? OnReconnected; public event Action? OnReconnected;
/// <inheritdoc />
public Func<Task<Uri?>>? GetReconnectionUrl { get; set; }
/// <summary> /// <summary>
/// ctor /// ctor
@@ -124,7 +118,7 @@ namespace CryptoExchange.Net.Sockets
Id = NextStreamId(); Id = NextStreamId();
_log = log; _log = log;
Parameters = websocketParameters; _parameters = websocketParameters;
_outgoingMessages = new List<DateTime>(); _outgoingMessages = new List<DateTime>();
_receivedMessages = new List<ReceiveItem>(); _receivedMessages = new List<ReceiveItem>();
_sendEvent = new AsyncResetEvent(); _sendEvent = new AsyncResetEvent();
@@ -153,26 +147,17 @@ namespace CryptoExchange.Net.Sockets
private ClientWebSocket CreateSocket() private ClientWebSocket CreateSocket()
{ {
var cookieContainer = new CookieContainer(); var cookieContainer = new CookieContainer();
foreach (var cookie in Parameters.Cookies) foreach (var cookie in _parameters.Cookies)
cookieContainer.Add(new Cookie(cookie.Key, cookie.Value)); cookieContainer.Add(new Cookie(cookie.Key, cookie.Value));
var socket = new ClientWebSocket(); var socket = new ClientWebSocket();
try socket.Options.Cookies = cookieContainer;
{ foreach (var header in _parameters.Headers)
socket.Options.Cookies = cookieContainer; socket.Options.SetRequestHeader(header.Key, header.Value);
foreach (var header in Parameters.Headers) socket.Options.KeepAliveInterval = _parameters.KeepAliveInterval ?? TimeSpan.Zero;
socket.Options.SetRequestHeader(header.Key, header.Value); socket.Options.SetBuffer(65536, 65536); // Setting it to anything bigger than 65536 throws an exception in .net framework
socket.Options.KeepAliveInterval = Parameters.KeepAliveInterval ?? TimeSpan.Zero; if (_parameters.Proxy != null)
socket.Options.SetBuffer(65536, 65536); // Setting it to anything bigger than 65536 throws an exception in .net framework SetProxy(_parameters.Proxy);
if (Parameters.Proxy != null)
SetProxy(Parameters.Proxy);
}
catch (PlatformNotSupportedException)
{
// Options are not supported on certain platforms (WebAssembly for instance)
// best we can do it try to connect without setting options.
}
return socket; return socket;
} }
@@ -203,7 +188,7 @@ namespace CryptoExchange.Net.Sockets
_processState = ProcessState.Processing; _processState = ProcessState.Processing;
var sendTask = SendLoopAsync(); var sendTask = SendLoopAsync();
var receiveTask = ReceiveLoopAsync(); var receiveTask = ReceiveLoopAsync();
var timeoutTask = Parameters.Timeout != null && Parameters.Timeout > TimeSpan.FromSeconds(0) ? CheckTimeoutAsync() : Task.CompletedTask; var timeoutTask = _parameters.Timeout != null && _parameters.Timeout > TimeSpan.FromSeconds(0) ? CheckTimeoutAsync() : Task.CompletedTask;
await Task.WhenAll(sendTask, receiveTask, timeoutTask).ConfigureAwait(false); await Task.WhenAll(sendTask, receiveTask, timeoutTask).ConfigureAwait(false);
_log.Write(LogLevel.Debug, $"Socket {Id} processing tasks finished"); _log.Write(LogLevel.Debug, $"Socket {Id} processing tasks finished");
@@ -214,7 +199,7 @@ namespace CryptoExchange.Net.Sockets
await _closeTask.ConfigureAwait(false); await _closeTask.ConfigureAwait(false);
_closeTask = null; _closeTask = null;
if (!Parameters.AutoReconnect) if (!_parameters.AutoReconnect)
{ {
_processState = ProcessState.Idle; _processState = ProcessState.Idle;
OnClose?.Invoke(); OnClose?.Invoke();
@@ -227,24 +212,9 @@ namespace CryptoExchange.Net.Sockets
OnReconnecting?.Invoke(); OnReconnecting?.Invoke();
} }
var sinceLastReconnect = DateTime.UtcNow - _lastReconnectTime;
if (sinceLastReconnect < Parameters.ReconnectInterval)
await Task.Delay(Parameters.ReconnectInterval - sinceLastReconnect).ConfigureAwait(false);
while (!_stopRequested) while (!_stopRequested)
{ {
_log.Write(LogLevel.Debug, $"Socket {Id} attempting to reconnect"); _log.Write(LogLevel.Debug, $"Socket {Id} attempting to reconnect");
var task = GetReconnectionUrl?.Invoke();
if (task != null)
{
var reconnectUri = await task.ConfigureAwait(false);
if (reconnectUri != null && Parameters.Uri != reconnectUri)
{
_log.Write(LogLevel.Debug, $"Socket {Id} reconnect URI set to {reconnectUri}");
Parameters.Uri = reconnectUri;
}
}
_socket = CreateSocket(); _socket = CreateSocket();
_ctsSource.Dispose(); _ctsSource.Dispose();
_ctsSource = new CancellationTokenSource(); _ctsSource = new CancellationTokenSource();
@@ -253,11 +223,10 @@ namespace CryptoExchange.Net.Sockets
var connected = await ConnectInternalAsync().ConfigureAwait(false); var connected = await ConnectInternalAsync().ConfigureAwait(false);
if (!connected) if (!connected)
{ {
await Task.Delay(Parameters.ReconnectInterval).ConfigureAwait(false); await Task.Delay(_parameters.ReconnectInterval).ConfigureAwait(false);
continue; continue;
} }
_lastReconnectTime = DateTime.UtcNow;
OnReconnected?.Invoke(); OnReconnected?.Invoke();
break; break;
} }
@@ -272,7 +241,7 @@ namespace CryptoExchange.Net.Sockets
if (_ctsSource.IsCancellationRequested) if (_ctsSource.IsCancellationRequested)
return; return;
var bytes = Parameters.Encoding.GetBytes(data); var bytes = _parameters.Encoding.GetBytes(data);
_log.Write(LogLevel.Trace, $"Socket {Id} Adding {bytes.Length} to sent buffer"); _log.Write(LogLevel.Trace, $"Socket {Id} Adding {bytes.Length} to sent buffer");
_sendBuffer.Enqueue(bytes); _sendBuffer.Enqueue(bytes);
_sendEvent.Set(); _sendEvent.Set();
@@ -281,7 +250,7 @@ namespace CryptoExchange.Net.Sockets
/// <inheritdoc /> /// <inheritdoc />
public virtual async Task ReconnectAsync() public virtual async Task ReconnectAsync()
{ {
if (_processState != ProcessState.Processing && IsOpen) if (_processState != ProcessState.Processing)
return; return;
_log.Write(LogLevel.Debug, $"Socket {Id} reconnect requested"); _log.Write(LogLevel.Debug, $"Socket {Id} reconnect requested");
@@ -401,11 +370,11 @@ namespace CryptoExchange.Net.Sockets
while (_sendBuffer.TryDequeue(out var data)) while (_sendBuffer.TryDequeue(out var data))
{ {
if (Parameters.RatelimitPerSecond != null) if (_parameters.RatelimitPerSecond != null)
{ {
// Wait for rate limit // Wait for rate limit
DateTime? start = null; DateTime? start = null;
while (MessagesSentLastSecond() >= Parameters.RatelimitPerSecond) while (MessagesSentLastSecond() >= _parameters.RatelimitPerSecond)
{ {
start ??= DateTime.UtcNow; start ??= DateTime.UtcNow;
await Task.Delay(50).ConfigureAwait(false); await Task.Delay(50).ConfigureAwait(false);
@@ -580,14 +549,14 @@ namespace CryptoExchange.Net.Sockets
string strData; string strData;
if (messageType == WebSocketMessageType.Binary) if (messageType == WebSocketMessageType.Binary)
{ {
if (Parameters.DataInterpreterBytes == null) if (_parameters.DataInterpreterBytes == null)
throw new Exception("Byte interpreter not set while receiving byte data"); throw new Exception("Byte interpreter not set while receiving byte data");
try try
{ {
var relevantData = new byte[count]; var relevantData = new byte[count];
Array.Copy(data, offset, relevantData, 0, count); Array.Copy(data, offset, relevantData, 0, count);
strData = Parameters.DataInterpreterBytes(relevantData); strData = _parameters.DataInterpreterBytes(relevantData);
} }
catch(Exception e) catch(Exception e)
{ {
@@ -596,13 +565,13 @@ namespace CryptoExchange.Net.Sockets
} }
} }
else else
strData = Parameters.Encoding.GetString(data, offset, count); strData = _parameters.Encoding.GetString(data, offset, count);
if (Parameters.DataInterpreterString != null) if (_parameters.DataInterpreterString != null)
{ {
try try
{ {
strData = Parameters.DataInterpreterString(strData); strData = _parameters.DataInterpreterString(strData);
} }
catch(Exception e) catch(Exception e)
{ {
@@ -664,7 +633,7 @@ namespace CryptoExchange.Net.Sockets
/// <returns></returns> /// <returns></returns>
protected async Task CheckTimeoutAsync() protected async Task CheckTimeoutAsync()
{ {
_log.Write(LogLevel.Debug, $"Socket {Id} Starting task checking for no data received for {Parameters.Timeout}"); _log.Write(LogLevel.Debug, $"Socket {Id} Starting task checking for no data received for {_parameters.Timeout}");
LastActionTime = DateTime.UtcNow; LastActionTime = DateTime.UtcNow;
try try
{ {
@@ -673,10 +642,10 @@ namespace CryptoExchange.Net.Sockets
if (_ctsSource.IsCancellationRequested) if (_ctsSource.IsCancellationRequested)
return; return;
if (DateTime.UtcNow - LastActionTime > Parameters.Timeout) if (DateTime.UtcNow - LastActionTime > _parameters.Timeout)
{ {
_log.Write(LogLevel.Warning, $"Socket {Id} No data received for {Parameters.Timeout}, reconnecting socket"); _log.Write(LogLevel.Warning, $"Socket {Id} No data received for {_parameters.Timeout}, reconnecting socket");
_ = ReconnectAsync().ConfigureAwait(false); _ = CloseAsync().ConfigureAwait(false);
return; return;
} }
try try
+1 -6
View File
@@ -25,12 +25,7 @@ namespace CryptoExchange.Net.Sockets
/// </summary> /// </summary>
public T Data { get; set; } public T Data { get; set; }
/// <summary> internal DataEvent(T data, DateTime timestamp)
/// Ctor
/// </summary>
/// <param name="data"></param>
/// <param name="timestamp"></param>
public DataEvent(T data, DateTime timestamp)
{ {
Data = data; Data = data;
Timestamp = timestamp; Timestamp = timestamp;
+3 -17
View File
@@ -184,7 +184,6 @@ namespace CryptoExchange.Net.Sockets
_socket.OnReconnecting += HandleReconnecting; _socket.OnReconnecting += HandleReconnecting;
_socket.OnReconnected += HandleReconnected; _socket.OnReconnected += HandleReconnected;
_socket.OnError += HandleError; _socket.OnError += HandleError;
_socket.GetReconnectionUrl = GetReconnectionUrlAsync;
} }
/// <summary> /// <summary>
@@ -224,17 +223,7 @@ namespace CryptoExchange.Net.Sockets
foreach (var sub in subscriptions) foreach (var sub in subscriptions)
sub.Confirmed = false; sub.Confirmed = false;
} }
Task.Run(() => ConnectionLost?.Invoke());
_ = Task.Run(() => ConnectionLost?.Invoke());
}
/// <summary>
/// Get the url to connect to when reconnecting
/// </summary>
/// <returns></returns>
protected virtual async Task<Uri?> GetReconnectionUrlAsync()
{
return await socketClient.GetReconnectUriAsync(ApiClient, this).ConfigureAwait(false);
} }
/// <summary> /// <summary>
@@ -396,7 +385,7 @@ namespace CryptoExchange.Net.Sockets
if (!subscriptions.Contains(subscription)) if (!subscriptions.Contains(subscription))
return; return;
subscription.Closed = true; subscriptions.Remove(subscription);
} }
if (Status == SocketStatus.Closing || Status == SocketStatus.Closed || Status == SocketStatus.Disposed) if (Status == SocketStatus.Closing || Status == SocketStatus.Closed || Status == SocketStatus.Disposed)
@@ -418,7 +407,7 @@ namespace CryptoExchange.Net.Sockets
return; return;
} }
shouldCloseConnection = subscriptions.All(r => !r.UserSubscription || r.Closed); shouldCloseConnection = subscriptions.All(r => !r.UserSubscription);
if (shouldCloseConnection) if (shouldCloseConnection)
Status = SocketStatus.Closing; Status = SocketStatus.Closing;
} }
@@ -428,9 +417,6 @@ namespace CryptoExchange.Net.Sockets
log.Write(LogLevel.Debug, $"Socket {SocketId} closing as there are no more subscriptions"); log.Write(LogLevel.Debug, $"Socket {SocketId} closing as there are no more subscriptions");
await CloseAsync().ConfigureAwait(false); await CloseAsync().ConfigureAwait(false);
} }
lock (subscriptionLock)
subscriptions.Remove(subscription);
} }
/// <summary> /// <summary>
@@ -48,11 +48,6 @@ namespace CryptoExchange.Net.Sockets
/// </summary> /// </summary>
public bool Authenticated { get; set; } public bool Authenticated { get; set; }
/// <summary>
/// Whether we're closing this subscription and a socket connection shouldn't be kept open for it
/// </summary>
public bool Closed { get; set; }
/// <summary> /// <summary>
/// Cancellation token registration, should be disposed when subscription is closed. Used for closing the subscription with /// Cancellation token registration, should be disposed when subscription is closed. Used for closing the subscription with
/// a provided cancelation token /// a provided cancelation token
@@ -2,7 +2,6 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Text; using System.Text;
using System.Threading.Tasks;
namespace CryptoExchange.Net.Sockets namespace CryptoExchange.Net.Sockets
{ {
+3 -35
View File
@@ -8,48 +8,16 @@ CryptoExchange.Net is a base package which can be used to easily implement crypt
## Discord ## Discord
A Discord server is available [here](https://discord.gg/MSpeEtSY8t). Feel free to join for discussion and/or questions around the CryptoExchange.Net and implementation libraries. A Discord server is available [here](https://discord.gg/MSpeEtSY8t). Feel free to join for discussion and/or questions around the CryptoExchange.Net and implementation libraries.
## Support the project ## Donate / Sponsor
I develop and maintain this package on my own for free in my spare time, any support is greatly appreciated. I develop and maintain this package on my own for free in my spare time. Donations are greatly appreciated. If you prefer to donate any other currency please contact me.
### Referral link
Use one of the following following referral links to signup to a new exchange to pay a small percentage of the trading fees you pay to support the project instead of paying them straight to the exchange. This doesn't cost you a thing!
[Binance](https://accounts.binance.com/en/register?ref=10153680)
[Bitfinex](https://www.bitfinex.com/sign-up?refcode=kCCe-CNBO)
[Bittrex](https://bittrex.com/discover/join?referralCode=TST-DJM-CSX)
[Bybit](https://partner.bybit.com/b/jkorf)
[CoinEx](https://www.coinex.com/register?refer_code=hd6gn)
[FTX](https://ftx.com/referrals#a=31620192)
[Huobi](https://www.huobi.com/en-us/v/register/double-invite/?inviter_id=11343840&invite_code=fxp93)
[Kucoin](https://www.kucoin.com/ucenter/signup?rcode=RguMux)
### Donate
Make a one time donation in a crypto currency of your choice. If you prefer to donate a currency not listed here please contact me.
**Btc**: 12KwZk3r2Y3JZ2uMULcjqqBvXmpDwjhhQS **Btc**: 12KwZk3r2Y3JZ2uMULcjqqBvXmpDwjhhQS
**Eth**: 0x069176ca1a4b1d6e0b7901a6bc0dbf3bb0bf5cc2 **Eth**: 0x069176ca1a4b1d6e0b7901a6bc0dbf3bb0bf5cc2
**Nano**: xrb_1ocs3hbp561ef76eoctjwg85w5ugr8wgimkj8mfhoyqbx4s1pbc74zggw7gs **Nano**: xrb_1ocs3hbp561ef76eoctjwg85w5ugr8wgimkj8mfhoyqbx4s1pbc74zggw7gs
### Sponsor 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.2.4 - 31 Jul 2022
* Added handling of PlatformNotSupportedException when trying to use websocket from WebAssembly
* Changed DataEvent to have a public constructor for testing purposes
* Fixed EnumConverter serializing values without proper quotes
* Fixed websocket connection reconnecting too quickly when resubscribing/reauthenticating fails
* Version 5.2.3 - 19 Jul 2022
* Fixed socket getting disconnected when `no data` timeout is reached instead of being reconnected
* Version 5.2.2 - 17 Jul 2022
* Added support for retrieving a new url when socket connection is lost and reconnection will happen
* Version 5.2.1 - 16 Jul 2022
* Fixed socket reconnect issue
* Fixed `message not handled` messages after unsubscribing
* Fixed error returning for non-json error responses
* Version 5.2.0 - 10 Jul 2022 * Version 5.2.0 - 10 Jul 2022
* Refactored websocket code, removed some clutter and simplified * Refactored websocket code, removed some clutter and simplified
* Added ReconnectAsync and GetSubscriptionsState methods on socket clients * Added ReconnectAsync and GetSubscriptionsState methods on socket clients
+1 -1
View File
@@ -105,7 +105,7 @@ All updates are wrapped in a `DataEvent<>` object, which contain a `Timestamp`,
*[WARNING] Do not use `using` statements in combination with constructing a `SocketClient`. Doing so will dispose the `SocketClient` instance when the subscription is done, which will result in the connection getting closed. Instead assign the socket client to a variable outside of the method scope.* *[WARNING] Do not use `using` statements in combination with constructing a `SocketClient`. Doing so will dispose the `SocketClient` instance when the subscription is done, which will result in the connection getting closed. Instead assign the socket client to a variable outside of the method scope.*
### Processing subscribe responses ### Processing subscribe responses
Subscribing to a stream will return a `CallResult<UpdateSubscription>` object. This should be checked for success the same way as the [rest client](#processing-request-responses). The `UpdateSubscription` object can be used to listen for connection events of the socket connection. Subscribing to a stream will return a `CallResult<UpdateSubscription>` object. This should be checked for success the same was as the [rest client](#processing-request-responses). The `UpdateSubscription` object can be used to listen for connection events of the socket connection.
```csharp ```csharp
var subscriptionResult = await kucoinSocketClient.SpotStreams.SubscribeToAllTickerUpdatesAsync(DataHandler); var subscriptionResult = await kucoinSocketClient.SpotStreams.SubscribeToAllTickerUpdatesAsync(DataHandler);
+3 -18
View File
@@ -42,26 +42,11 @@ These might not be compatible with other libraries, make sure to check the Crypt
## Discord ## Discord
A Discord server is available [here](https://discord.gg/MSpeEtSY8t). Feel free to join for discussion and/or questions around the CryptoExchange.Net and implementation libraries. A Discord server is available [here](https://discord.gg/MSpeEtSY8t). Feel free to join for discussion and/or questions around the CryptoExchange.Net and implementation libraries.
## Support the project ## Donate / Sponsor
I develop and maintain this package on my own for free in my spare time, any support is greatly appreciated. I develop and maintain this package on my own for free in my spare time. Donations are greatly appreciated. If you prefer to donate any other currency please contact me.
### Referral link
Use one of the following following referral links to signup to a new exchange to pay a small percentage of the trading fees you pay to support the project instead of paying them straight to the exchange. This doesn't cost you a thing!
[Binance](https://accounts.binance.com/en/register?ref=10153680)
[Bitfinex](https://www.bitfinex.com/sign-up?refcode=kCCe-CNBO)
[Bittrex](https://bittrex.com/discover/join?referralCode=TST-DJM-CSX)
[Bybit](https://partner.bybit.com/b/jkorf)
[CoinEx](https://www.coinex.com/register?refer_code=hd6gn)
[FTX](https://ftx.com/referrals#a=31620192)
[Huobi](https://www.huobi.com/en-us/v/register/double-invite/?inviter_id=11343840&invite_code=fxp93)
[Kucoin](https://www.kucoin.com/ucenter/signup?rcode=RguMux)
### Donate
Make a one time donation in a crypto currency of your choice. If you prefer to donate a currency not listed here please contact me.
**Btc**: 12KwZk3r2Y3JZ2uMULcjqqBvXmpDwjhhQS **Btc**: 12KwZk3r2Y3JZ2uMULcjqqBvXmpDwjhhQS
**Eth**: 0x069176ca1a4b1d6e0b7901a6bc0dbf3bb0bf5cc2 **Eth**: 0x069176ca1a4b1d6e0b7901a6bc0dbf3bb0bf5cc2
**Nano**: xrb_1ocs3hbp561ef76eoctjwg85w5ugr8wgimkj8mfhoyqbx4s1pbc74zggw7gs **Nano**: xrb_1ocs3hbp561ef76eoctjwg85w5ugr8wgimkj8mfhoyqbx4s1pbc74zggw7gs
### Sponsor 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).