mirror of
https://github.com/JKorf/CryptoExchange.Net.git
synced 2026-08-11 08:22:53 +00:00
Added ManualUpdateSubscription and UpdateSubscription additional constructor to allow producing websocket events without actual connection
This commit is contained in:
@@ -0,0 +1,136 @@
|
||||
using CryptoExchange.Net.Objects;
|
||||
using CryptoExchange.Net.Objects.Errors;
|
||||
using CryptoExchange.Net.Objects.Sockets;
|
||||
using CryptoExchange.Net.Sockets.Default;
|
||||
using NUnit.Framework;
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace CryptoExchange.Net.UnitTests
|
||||
{
|
||||
[TestFixture]
|
||||
public class ManualUpdateSubscriptionTests
|
||||
{
|
||||
[Test]
|
||||
public void Constructor_Should_CreateSubscribedVirtualSubscription()
|
||||
{
|
||||
var controller = new ManualUpdateSubscription(socketId: 12);
|
||||
|
||||
Assert.That(controller.Subscription.SocketId, Is.EqualTo(12));
|
||||
Assert.That(controller.Subscription.Id, Is.GreaterThan(0));
|
||||
Assert.That(controller.Subscription.SocketStatus, Is.EqualTo(SocketStatus.Connected));
|
||||
Assert.That(controller.Subscription.SubscriptionStatus, Is.EqualTo(SubscriptionStatus.Subscribed));
|
||||
Assert.That(controller.Subscription.LastReceiveTime, Is.Null);
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void StateChanges_Should_BeVisibleOnSubscription()
|
||||
{
|
||||
var controller = new ManualUpdateSubscription();
|
||||
var timestamp = new DateTime(2026, 8, 5, 12, 0, 0, DateTimeKind.Utc);
|
||||
var statuses = new List<SubscriptionStatus>();
|
||||
controller.Subscription.SubscriptionStatusChanged += statuses.Add;
|
||||
|
||||
controller.SetLastReceiveTime(timestamp);
|
||||
controller.SetSocketStatus(SocketStatus.Reconnecting);
|
||||
controller.SetSubscriptionStatus(SubscriptionStatus.Subscribing);
|
||||
controller.SetSubscriptionStatus(SubscriptionStatus.Subscribed);
|
||||
|
||||
Assert.That(controller.Subscription.LastReceiveTime, Is.EqualTo(timestamp));
|
||||
Assert.That(controller.Subscription.SocketStatus, Is.EqualTo(SocketStatus.Reconnecting));
|
||||
Assert.That(controller.Subscription.SubscriptionStatus, Is.EqualTo(SubscriptionStatus.Subscribed));
|
||||
Assert.That(statuses, Is.EqualTo(new[]
|
||||
{
|
||||
SubscriptionStatus.Subscribing,
|
||||
SubscriptionStatus.Subscribed
|
||||
}));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void LifecycleMethods_Should_InvokeSubscriptionEvents()
|
||||
{
|
||||
var controller = new ManualUpdateSubscription();
|
||||
var error = new ServerError("Test error", ErrorInfo.Unknown);
|
||||
var exception = new InvalidOperationException("Test exception");
|
||||
var disconnectedPeriod = TimeSpan.FromMinutes(2);
|
||||
var lost = 0;
|
||||
var restored = TimeSpan.Zero;
|
||||
Error? resubscribeError = null;
|
||||
var paused = 0;
|
||||
var unpaused = 0;
|
||||
Exception? receivedException = null;
|
||||
|
||||
controller.Subscription.ConnectionLost += () => lost++;
|
||||
controller.Subscription.ConnectionRestored += x => restored = x;
|
||||
controller.Subscription.ResubscribingFailed += x => resubscribeError = x;
|
||||
controller.Subscription.ActivityPaused += () => paused++;
|
||||
controller.Subscription.ActivityUnpaused += () => unpaused++;
|
||||
controller.Subscription.Exception += x => receivedException = x;
|
||||
|
||||
controller.InvokeConnectionLost();
|
||||
controller.InvokeConnectionRestored(disconnectedPeriod);
|
||||
controller.InvokeResubscribingFailed(error);
|
||||
controller.InvokeActivityPaused();
|
||||
controller.InvokeActivityUnpaused();
|
||||
controller.InvokeException(exception);
|
||||
|
||||
Assert.That(lost, Is.EqualTo(1));
|
||||
Assert.That(restored, Is.EqualTo(disconnectedPeriod));
|
||||
Assert.That(resubscribeError, Is.SameAs(error));
|
||||
Assert.That(paused, Is.EqualTo(1));
|
||||
Assert.That(unpaused, Is.EqualTo(1));
|
||||
Assert.That(receivedException, Is.SameAs(exception));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void InvokeConnectionClosed_Should_CloseAndOnlyInvokeOnce()
|
||||
{
|
||||
var controller = new ManualUpdateSubscription();
|
||||
var closed = 0;
|
||||
controller.Subscription.ConnectionClosed += () => closed++;
|
||||
|
||||
controller.InvokeConnectionClosed();
|
||||
controller.InvokeConnectionClosed();
|
||||
|
||||
Assert.That(closed, Is.EqualTo(1));
|
||||
Assert.That(controller.Subscription.SocketStatus, Is.EqualTo(SocketStatus.Closed));
|
||||
Assert.That(controller.Subscription.SubscriptionStatus, Is.EqualTo(SubscriptionStatus.Closed));
|
||||
}
|
||||
|
||||
[Test]
|
||||
public async Task SubscriptionOperations_Should_InvokeCallbacks()
|
||||
{
|
||||
var closes = 0;
|
||||
var reconnects = 0;
|
||||
var resubscribes = 0;
|
||||
var controller = new ManualUpdateSubscription(
|
||||
closeAsync: () =>
|
||||
{
|
||||
closes++;
|
||||
return Task.CompletedTask;
|
||||
},
|
||||
reconnectAsync: () =>
|
||||
{
|
||||
reconnects++;
|
||||
return Task.CompletedTask;
|
||||
},
|
||||
resubscribeAsync: () =>
|
||||
{
|
||||
resubscribes++;
|
||||
return Task.FromResult(CallResult.Ok());
|
||||
});
|
||||
|
||||
await controller.Subscription.ReconnectAsync();
|
||||
var resubscribeResult = await controller.Subscription.ResubscribeAsync();
|
||||
await controller.Subscription.CloseAsync();
|
||||
await controller.Subscription.CloseAsync();
|
||||
|
||||
Assert.That(reconnects, Is.EqualTo(1));
|
||||
Assert.That(resubscribes, Is.EqualTo(1));
|
||||
Assert.That(resubscribeResult.Success, Is.True);
|
||||
Assert.That(closes, Is.EqualTo(1));
|
||||
Assert.That(controller.Subscription.SubscriptionStatus, Is.EqualTo(SubscriptionStatus.Closed));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,209 @@
|
||||
using CryptoExchange.Net.Sockets;
|
||||
using CryptoExchange.Net.Sockets.Default;
|
||||
using CryptoExchange.Net.Sockets.Default.Routing;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace CryptoExchange.Net.Objects.Sockets
|
||||
{
|
||||
/// <summary>
|
||||
/// Controller for an update subscription which isn't backed by a websocket connection. Can be used for testing.
|
||||
/// </summary>
|
||||
public class ManualUpdateSubscription
|
||||
{
|
||||
private readonly Func<Task> _closeAsync;
|
||||
private readonly Func<Task> _reconnectAsync;
|
||||
private readonly Func<Task<CallResult>> _resubscribeAsync;
|
||||
private readonly ManualSubscription _manualSubscription;
|
||||
private int _closedEventInvoked;
|
||||
|
||||
/// <summary>
|
||||
/// The update subscription
|
||||
/// </summary>
|
||||
public UpdateSubscription Subscription { get; }
|
||||
|
||||
/// <summary>
|
||||
/// The virtual socket id
|
||||
/// </summary>
|
||||
public int SocketId { get; }
|
||||
|
||||
/// <summary>
|
||||
/// The last timestamp anything was received by the subscription
|
||||
/// </summary>
|
||||
public DateTime? LastReceiveTime { get; private set; }
|
||||
|
||||
/// <summary>
|
||||
/// The current virtual websocket status
|
||||
/// </summary>
|
||||
public SocketStatus SocketStatus { get; private set; }
|
||||
|
||||
/// <summary>
|
||||
/// Create a manually controlled update subscription
|
||||
/// </summary>
|
||||
/// <param name="socketId">The virtual socket id</param>
|
||||
/// <param name="closeAsync">Callback when the subscription is closed</param>
|
||||
/// <param name="reconnectAsync">Callback when a reconnect is requested</param>
|
||||
/// <param name="resubscribeAsync">Callback when a resubscribe is requested</param>
|
||||
public ManualUpdateSubscription(
|
||||
int socketId = 0,
|
||||
Func<Task>? closeAsync = null,
|
||||
Func<Task>? reconnectAsync = null,
|
||||
Func<Task<CallResult>>? resubscribeAsync = null)
|
||||
{
|
||||
SocketId = socketId;
|
||||
SocketStatus = SocketStatus.Connected;
|
||||
_closeAsync = closeAsync ?? (() => Task.CompletedTask);
|
||||
_reconnectAsync = reconnectAsync ?? (() => Task.CompletedTask);
|
||||
_resubscribeAsync = resubscribeAsync ?? (() => Task.FromResult(CallResult.Ok()));
|
||||
|
||||
_manualSubscription = new ManualSubscription();
|
||||
_manualSubscription.Status = SubscriptionStatus.Subscribed;
|
||||
Subscription = new UpdateSubscription(this, _manualSubscription);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Set the last timestamp anything was received by the subscription
|
||||
/// </summary>
|
||||
/// <param name="timestamp">The receive timestamp</param>
|
||||
public void SetLastReceiveTime(DateTime? timestamp)
|
||||
{
|
||||
LastReceiveTime = timestamp;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Set the virtual websocket status
|
||||
/// </summary>
|
||||
/// <param name="status">The status</param>
|
||||
public void SetSocketStatus(SocketStatus status)
|
||||
{
|
||||
SocketStatus = status;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Set the subscription status
|
||||
/// </summary>
|
||||
/// <param name="status">The status</param>
|
||||
public void SetSubscriptionStatus(SubscriptionStatus status)
|
||||
{
|
||||
_manualSubscription.Status = status;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the connection lost event
|
||||
/// </summary>
|
||||
public void InvokeConnectionLost()
|
||||
{
|
||||
Subscription.HandleConnectionLostEvent();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the connection restored event
|
||||
/// </summary>
|
||||
/// <param name="disconnectedPeriod">The period the connection was disconnected</param>
|
||||
public void InvokeConnectionRestored(TimeSpan disconnectedPeriod)
|
||||
{
|
||||
Subscription.HandleConnectionRestoredEvent(disconnectedPeriod);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the connection closed event
|
||||
/// </summary>
|
||||
public void InvokeConnectionClosed()
|
||||
{
|
||||
if (Interlocked.Exchange(ref _closedEventInvoked, 1) != 0)
|
||||
return;
|
||||
|
||||
SocketStatus = SocketStatus.Closed;
|
||||
_manualSubscription.Status = SubscriptionStatus.Closed;
|
||||
Subscription.HandleConnectionClosedEvent();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the resubscribing failed event
|
||||
/// </summary>
|
||||
/// <param name="error">The resubscribe error</param>
|
||||
public void InvokeResubscribingFailed(Error error)
|
||||
{
|
||||
if (error == null)
|
||||
throw new ArgumentNullException(nameof(error));
|
||||
|
||||
Subscription.HandleResubscribeFailedEvent(error);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the activity paused event
|
||||
/// </summary>
|
||||
public void InvokeActivityPaused()
|
||||
{
|
||||
Subscription.HandlePausedEvent();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the activity unpaused event
|
||||
/// </summary>
|
||||
public void InvokeActivityUnpaused()
|
||||
{
|
||||
Subscription.HandleUnpausedEvent();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Invoke the exception event
|
||||
/// </summary>
|
||||
/// <param name="exception">The exception</param>
|
||||
public void InvokeException(Exception exception)
|
||||
{
|
||||
if (exception == null)
|
||||
throw new ArgumentNullException(nameof(exception));
|
||||
|
||||
_manualSubscription.InvokeExceptionHandler(exception);
|
||||
}
|
||||
|
||||
internal async Task CloseAsync()
|
||||
{
|
||||
if (_manualSubscription.Status == SubscriptionStatus.Closed
|
||||
|| _manualSubscription.Status == SubscriptionStatus.Closing)
|
||||
return;
|
||||
|
||||
_manualSubscription.Status = SubscriptionStatus.Closing;
|
||||
try
|
||||
{
|
||||
await _closeAsync().ConfigureAwait(false);
|
||||
}
|
||||
finally
|
||||
{
|
||||
_manualSubscription.Status = SubscriptionStatus.Closed;
|
||||
}
|
||||
}
|
||||
|
||||
internal Task ReconnectAsync()
|
||||
{
|
||||
return _reconnectAsync();
|
||||
}
|
||||
|
||||
internal Task<CallResult> ResubscribeAsync()
|
||||
{
|
||||
return _resubscribeAsync();
|
||||
}
|
||||
|
||||
private class ManualSubscription : Subscription
|
||||
{
|
||||
public ManualSubscription()
|
||||
: base(NullLogger.Instance, false)
|
||||
{
|
||||
MessageRouter = MessageRouter.Create();
|
||||
}
|
||||
|
||||
protected override Query? GetSubQuery(SocketConnection connection)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
protected override Query? GetUnsubQuery(SocketConnection connection)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -12,7 +12,8 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// </summary>
|
||||
public class UpdateSubscription
|
||||
{
|
||||
private readonly SocketConnection _connection;
|
||||
private readonly SocketConnection? _connection;
|
||||
private readonly ManualUpdateSubscription? _manualSubscription;
|
||||
internal readonly Subscription _subscription;
|
||||
|
||||
#if NET9_0_OR_GREATER
|
||||
@@ -102,7 +103,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <summary>
|
||||
/// The id of the socket
|
||||
/// </summary>
|
||||
public int SocketId => _connection.SocketId;
|
||||
public int SocketId => _connection?.SocketId ?? _manualSubscription!.SocketId;
|
||||
|
||||
/// <summary>
|
||||
/// The id of the subscription
|
||||
@@ -112,12 +113,12 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <summary>
|
||||
/// The last timestamp anything was received from the server
|
||||
/// </summary>
|
||||
public DateTime? LastReceiveTime => _connection.LastReceiveTime;
|
||||
public DateTime? LastReceiveTime => _connection?.LastReceiveTime ?? _manualSubscription!.LastReceiveTime;
|
||||
|
||||
/// <summary>
|
||||
/// The current websocket status
|
||||
/// </summary>
|
||||
public SocketStatus SocketStatus => _connection.Status;
|
||||
public SocketStatus SocketStatus => _connection?.Status ?? _manualSubscription!.SocketStatus;
|
||||
|
||||
/// <summary>
|
||||
/// The current subscription status
|
||||
@@ -143,6 +144,18 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
_subscription.StatusChanged += (x) => SubscriptionStatusChanged?.Invoke(x);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// ctor
|
||||
/// </summary>
|
||||
/// <param name="manualSubscription">The manual subscription for controlling events and data</param>
|
||||
/// <param name="subscription">The subscription</param>
|
||||
internal UpdateSubscription(ManualUpdateSubscription manualSubscription, Subscription subscription)
|
||||
{
|
||||
_manualSubscription = manualSubscription;
|
||||
_subscription = subscription;
|
||||
_subscription.StatusChanged += (x) => SubscriptionStatusChanged?.Invoke(x);
|
||||
}
|
||||
|
||||
private void UnsubscribeConnectionEvents()
|
||||
{
|
||||
lock (_eventLock)
|
||||
@@ -150,22 +163,26 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
if (!_connectionEventsSubscribed)
|
||||
return;
|
||||
|
||||
_connection.ConnectionClosed -= HandleConnectionClosedEvent;
|
||||
_connection.ConnectionLost -= HandleConnectionLostEvent;
|
||||
_connection.ConnectionRestored -= HandleConnectionRestoredEvent;
|
||||
_connection.ResubscribingFailed -= HandleResubscribeFailedEvent;
|
||||
_connection.ActivityPaused -= HandlePausedEvent;
|
||||
_connection.ActivityUnpaused -= HandleUnpausedEvent;
|
||||
if (_connection != null)
|
||||
{
|
||||
_connection.ConnectionClosed -= HandleConnectionClosedEvent;
|
||||
_connection.ConnectionLost -= HandleConnectionLostEvent;
|
||||
_connection.ConnectionRestored -= HandleConnectionRestoredEvent;
|
||||
_connection.ResubscribingFailed -= HandleResubscribeFailedEvent;
|
||||
_connection.ActivityPaused -= HandlePausedEvent;
|
||||
_connection.ActivityUnpaused -= HandleUnpausedEvent;
|
||||
}
|
||||
|
||||
_connectionEventsSubscribed = false;
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleConnectionClosedEvent()
|
||||
internal void HandleConnectionClosedEvent()
|
||||
{
|
||||
UnsubscribeConnectionEvents();
|
||||
|
||||
// If we're not the subscription closing this connection don't bother emitting
|
||||
if (!_subscription.IsClosingConnection)
|
||||
if (_connection != null && !_subscription.IsClosingConnection)
|
||||
return;
|
||||
|
||||
List<Action> handlers;
|
||||
@@ -176,7 +193,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
callback();
|
||||
}
|
||||
|
||||
private void HandleConnectionLostEvent()
|
||||
internal void HandleConnectionLostEvent()
|
||||
{
|
||||
if (!_subscription.Active)
|
||||
{
|
||||
@@ -192,7 +209,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
callback();
|
||||
}
|
||||
|
||||
private void HandleConnectionRestoredEvent(TimeSpan period)
|
||||
internal void HandleConnectionRestoredEvent(TimeSpan period)
|
||||
{
|
||||
if (!_subscription.Active)
|
||||
{
|
||||
@@ -208,7 +225,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
callback(period);
|
||||
}
|
||||
|
||||
private void HandleResubscribeFailedEvent(Error error)
|
||||
internal void HandleResubscribeFailedEvent(Error error)
|
||||
{
|
||||
if (!_subscription.Active)
|
||||
{
|
||||
@@ -224,7 +241,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
callback(error);
|
||||
}
|
||||
|
||||
private void HandlePausedEvent()
|
||||
internal void HandlePausedEvent()
|
||||
{
|
||||
if (!_subscription.Active)
|
||||
{
|
||||
@@ -240,7 +257,7 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
callback();
|
||||
}
|
||||
|
||||
private void HandleUnpausedEvent()
|
||||
internal void HandleUnpausedEvent()
|
||||
{
|
||||
if (!_subscription.Active)
|
||||
{
|
||||
@@ -262,7 +279,10 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <returns></returns>
|
||||
public Task CloseAsync()
|
||||
{
|
||||
return _connection.CloseAsync(_subscription);
|
||||
if (_connection != null)
|
||||
return _connection.CloseAsync(_subscription);
|
||||
|
||||
return _manualSubscription!.CloseAsync();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -271,7 +291,10 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <returns></returns>
|
||||
public Task ReconnectAsync()
|
||||
{
|
||||
return _connection.TriggerReconnectAsync();
|
||||
if (_connection != null)
|
||||
return _connection.TriggerReconnectAsync();
|
||||
|
||||
return _manualSubscription!.ReconnectAsync();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -280,7 +303,13 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <returns></returns>
|
||||
internal async Task UnsubscribeAsync()
|
||||
{
|
||||
await _connection.UnsubscribeAsync(_subscription).ConfigureAwait(false);
|
||||
if (_connection != null)
|
||||
{
|
||||
await _connection.UnsubscribeAsync(_subscription).ConfigureAwait(false);
|
||||
return;
|
||||
}
|
||||
|
||||
await _manualSubscription!.CloseAsync().ConfigureAwait(false);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -289,7 +318,10 @@ namespace CryptoExchange.Net.Objects.Sockets
|
||||
/// <returns></returns>
|
||||
internal async Task<CallResult> ResubscribeAsync()
|
||||
{
|
||||
return await _connection.ResubscribeAsync(_subscription).ConfigureAwait(false);
|
||||
if (_connection != null)
|
||||
return await _connection.ResubscribeAsync(_subscription).ConfigureAwait(false);
|
||||
|
||||
return await _manualSubscription!.ResubscribeAsync().ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user