Add TestClient to simplify test code
This commit is contained in:
parent
5def499323
commit
72e7b54dff
|
|
@ -2,18 +2,11 @@
|
|||
// Licensed under the Apache License, Version 2.0. See License.txt in the project root for license information.
|
||||
|
||||
using System;
|
||||
using System.IO;
|
||||
using System.IO.Pipelines;
|
||||
using System.Security.Claims;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using System.Threading.Tasks.Channels;
|
||||
using Microsoft.AspNetCore.Sockets;
|
||||
using Microsoft.AspNetCore.Sockets.Internal;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Internal;
|
||||
using Moq;
|
||||
using Newtonsoft.Json;
|
||||
using Xunit;
|
||||
|
||||
namespace Microsoft.AspNetCore.SignalR.Tests
|
||||
|
|
@ -27,12 +20,12 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
var serviceProvider = CreateServiceProvider(s => s.AddSingleton(trackDispose));
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<TestHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask;
|
||||
|
||||
|
|
@ -57,14 +50,14 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<Hub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var exception =
|
||||
await Assert.ThrowsAsync<InvalidOperationException>(
|
||||
async () => await endPoint.OnConnectedAsync(connectionWrapper.Connection));
|
||||
async () => await endPoint.OnConnectedAsync(client.Connection));
|
||||
Assert.Equal("Lifetime manager OnConnectedAsync failed.", exception.Message);
|
||||
|
||||
connectionWrapper.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
mockLifetimeManager.Verify(m => m.OnConnectedAsync(It.IsAny<Connection>()), Times.Once);
|
||||
mockLifetimeManager.Verify(m => m.OnDisconnectedAsync(It.IsAny<Connection>()), Times.Once);
|
||||
|
|
@ -85,10 +78,10 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<OnConnectedThrowsHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
connectionWrapper.Dispose();
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
client.Dispose();
|
||||
|
||||
var exception = await Assert.ThrowsAsync<InvalidOperationException>(async () => await endPointTask);
|
||||
Assert.Equal("Hub OnConnected failed.", exception.Message);
|
||||
|
|
@ -109,10 +102,10 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<OnDisconnectedThrowsHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
connectionWrapper.Dispose();
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
client.Dispose();
|
||||
|
||||
var exception = await Assert.ThrowsAsync<InvalidOperationException>(async () => await endPointTask);
|
||||
Assert.Equal("Hub OnDisconnected failed.", exception.Message);
|
||||
|
|
@ -129,21 +122,17 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, nameof(MethodHub.TaskValueMethod));
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.TaskValueMethod)).OrTimeout();
|
||||
|
||||
// json serializer makes this a long
|
||||
Assert.Equal(42L, result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -156,21 +145,17 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, "echo", "hello");
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>("echo", "hello").OrTimeout();
|
||||
|
||||
Assert.Null(result.Error);
|
||||
Assert.Equal("hello", result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -183,21 +168,17 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, nameof(MethodHub.ValueMethod));
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.ValueMethod)).OrTimeout();
|
||||
|
||||
// json serializer makes this a long
|
||||
Assert.Equal(43L, result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -210,20 +191,16 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, "StaticMethod");
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.StaticMethod)).OrTimeout();
|
||||
|
||||
Assert.Equal("fromStatic", result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -236,20 +213,16 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, "VoidMethod");
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.VoidMethod)).OrTimeout();
|
||||
|
||||
Assert.Null(result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -262,21 +235,18 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.ConcatString), (byte)32, 42, 'm', "string").OrTimeout();
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, "ConcatString", (byte)32, 42, 'm', "string");
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
Assert.Equal("32, 42, m, string", result.Result);
|
||||
|
||||
// kill the connection
|
||||
connectionWrapper.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask;
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -287,17 +257,18 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connectionWrapper = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connectionWrapper.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(connectionWrapper, adapter, "OnDisconnectedAsync");
|
||||
var result = await ReadConnectionOutputAsync<InvocationResultDescriptor>(connectionWrapper).OrTimeout();
|
||||
var result = await client.Invoke<InvocationResultDescriptor>(nameof(MethodHub.OnDisconnectedAsync)).OrTimeout();
|
||||
|
||||
Assert.Equal("Unknown hub method 'OnDisconnectedAsync'", result.Error);
|
||||
|
||||
// kill the connection
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -308,22 +279,19 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var firstConnection = new ConnectionWrapper())
|
||||
using (var secondConnection = new ConnectionWrapper())
|
||||
using (var firstClient = new TestClient(serviceProvider))
|
||||
using (var secondClient = new TestClient(serviceProvider))
|
||||
{
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstConnection.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondConnection.Connection);
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstClient.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondClient.Connection);
|
||||
|
||||
await Task.WhenAll(firstConnection.Connected, secondConnection.Connected).OrTimeout();
|
||||
await Task.WhenAll(firstClient.Connected, secondClient.Connected).OrTimeout();
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest(firstConnection, adapter, "BroadcastMethod", "test");
|
||||
await firstClient.Invoke(nameof(MethodHub.BroadcastMethod), "test").OrTimeout();
|
||||
|
||||
foreach (var result in await Task.WhenAll(
|
||||
ReadConnectionOutputAsync<InvocationDescriptor>(firstConnection),
|
||||
ReadConnectionOutputAsync<InvocationDescriptor>(secondConnection)).OrTimeout())
|
||||
firstClient.Read<InvocationDescriptor>(),
|
||||
secondClient.Read<InvocationDescriptor>()).OrTimeout())
|
||||
{
|
||||
Assert.Equal("Broadcast", result.Method);
|
||||
Assert.Equal(1, result.Arguments.Length);
|
||||
|
|
@ -331,8 +299,8 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
}
|
||||
|
||||
// kill the connections
|
||||
firstConnection.Connection.Dispose();
|
||||
secondConnection.Connection.Dispose();
|
||||
firstClient.Dispose();
|
||||
secondClient.Dispose();
|
||||
|
||||
await Task.WhenAll(firstEndPointTask, secondEndPointTask).OrTimeout();
|
||||
}
|
||||
|
|
@ -345,38 +313,35 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var firstConnection = new ConnectionWrapper())
|
||||
using (var secondConnection = new ConnectionWrapper())
|
||||
using (var firstClient = new TestClient(serviceProvider))
|
||||
using (var secondClient = new TestClient(serviceProvider))
|
||||
{
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstConnection.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondConnection.Connection);
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstClient.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondClient.Connection);
|
||||
|
||||
await Task.WhenAll(firstConnection.Connected, secondConnection.Connected).OrTimeout();
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest_IgnoreReceive(firstConnection, adapter, "GroupSendMethod", "testGroup", "test");
|
||||
// check that 'secondConnection' hasn't received the group send
|
||||
Message message;
|
||||
Assert.False(secondConnection.Application.Input.TryRead(out message));
|
||||
|
||||
await SendRequest_IgnoreReceive(secondConnection, adapter, "GroupAddMethod", "testGroup");
|
||||
|
||||
await SendRequest(firstConnection, adapter, "GroupSendMethod", "testGroup", "test");
|
||||
await Task.WhenAll(firstClient.Connected, secondClient.Connected).OrTimeout();
|
||||
|
||||
var result = await firstClient.Invoke<InvocationResultDescriptor>(nameof(MethodHub.GroupSendMethod), "testGroup", "test").OrTimeout();
|
||||
// check that 'firstConnection' hasn't received the group send
|
||||
Assert.False(firstConnection.Application.Input.TryRead(out message));
|
||||
Assert.Null(result.Id);
|
||||
|
||||
// check that 'secondConnection' hasn't received the group send
|
||||
Assert.Null(await secondClient.TryRead<InvocationDescriptor>().OrTimeout());
|
||||
|
||||
result = await secondClient.Invoke<InvocationResultDescriptor>(nameof(MethodHub.GroupAddMethod), "testGroup").OrTimeout();
|
||||
Assert.Null(result.Id);
|
||||
|
||||
await firstClient.Invoke(nameof(MethodHub.GroupSendMethod), "testGroup", "test").OrTimeout();
|
||||
|
||||
// check that 'secondConnection' has received the group send
|
||||
var res = await ReadConnectionOutputAsync<InvocationDescriptor>(secondConnection).OrTimeout();
|
||||
Assert.Equal("Send", res.Method);
|
||||
Assert.Equal(1, res.Arguments.Length);
|
||||
Assert.Equal("test", res.Arguments[0]);
|
||||
var descriptor = await secondClient.Read<InvocationDescriptor>().OrTimeout();
|
||||
Assert.Equal("Send", descriptor.Method);
|
||||
Assert.Equal(1, descriptor.Arguments.Length);
|
||||
Assert.Equal("test", descriptor.Arguments[0]);
|
||||
|
||||
// kill the connections
|
||||
firstConnection.Connection.Dispose();
|
||||
secondConnection.Connection.Dispose();
|
||||
firstClient.Dispose();
|
||||
secondClient.Dispose();
|
||||
|
||||
await Task.WhenAll(firstEndPointTask, secondEndPointTask).OrTimeout();
|
||||
}
|
||||
|
|
@ -389,17 +354,14 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var connection = new ConnectionWrapper())
|
||||
using (var client = new TestClient(serviceProvider))
|
||||
{
|
||||
var endPointTask = endPoint.OnConnectedAsync(connection.Connection);
|
||||
var endPointTask = endPoint.OnConnectedAsync(client.Connection);
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var writer = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest_IgnoreReceive(connection, writer, "GroupRemoveMethod", "testGroup");
|
||||
await client.Invoke(nameof(MethodHub.GroupRemoveMethod), "testGroup").OrTimeout();
|
||||
|
||||
// kill the connection
|
||||
connection.Connection.Dispose();
|
||||
client.Dispose();
|
||||
|
||||
await endPointTask.OrTimeout();
|
||||
}
|
||||
|
|
@ -412,28 +374,25 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var firstConnection = new ConnectionWrapper())
|
||||
using (var secondConnection = new ConnectionWrapper())
|
||||
using (var firstClient = new TestClient(serviceProvider))
|
||||
using (var secondClient = new TestClient(serviceProvider))
|
||||
{
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstConnection.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondConnection.Connection);
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstClient.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondClient.Connection);
|
||||
|
||||
await Task.WhenAll(firstConnection.Connected, secondConnection.Connected).OrTimeout();
|
||||
await Task.WhenAll(firstClient.Connected, secondClient.Connected).OrTimeout();
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest_IgnoreReceive(firstConnection, adapter, "ClientSendMethod", secondConnection.Connection.User.Identity.Name, "test");
|
||||
await firstClient.Invoke(nameof(MethodHub.ClientSendMethod), secondClient.Connection.User.Identity.Name, "test").OrTimeout();
|
||||
|
||||
// check that 'secondConnection' has received the group send
|
||||
var res = await ReadConnectionOutputAsync<InvocationDescriptor>(secondConnection).OrTimeout();
|
||||
Assert.Equal("Send", res.Method);
|
||||
Assert.Equal(1, res.Arguments.Length);
|
||||
Assert.Equal("test", res.Arguments[0]);
|
||||
var result = await secondClient.Read<InvocationDescriptor>().OrTimeout();
|
||||
Assert.Equal("Send", result.Method);
|
||||
Assert.Equal(1, result.Arguments.Length);
|
||||
Assert.Equal("test", result.Arguments[0]);
|
||||
|
||||
// kill the connections
|
||||
firstConnection.Connection.Dispose();
|
||||
secondConnection.Connection.Dispose();
|
||||
firstClient.Dispose();
|
||||
secondClient.Dispose();
|
||||
|
||||
await Task.WhenAll(firstEndPointTask, secondEndPointTask).OrTimeout();
|
||||
}
|
||||
|
|
@ -446,28 +405,25 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
|
||||
var endPoint = serviceProvider.GetService<HubEndPoint<MethodHub>>();
|
||||
|
||||
using (var firstConnection = new ConnectionWrapper())
|
||||
using (var secondConnection = new ConnectionWrapper())
|
||||
using (var firstClient = new TestClient(serviceProvider))
|
||||
using (var secondClient = new TestClient(serviceProvider))
|
||||
{
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstConnection.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondConnection.Connection);
|
||||
var firstEndPointTask = endPoint.OnConnectedAsync(firstClient.Connection);
|
||||
var secondEndPointTask = endPoint.OnConnectedAsync(secondClient.Connection);
|
||||
|
||||
await Task.WhenAll(firstConnection.Connected, secondConnection.Connected).OrTimeout();
|
||||
await Task.WhenAll(firstClient.Connected, secondClient.Connected).OrTimeout();
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
var adapter = invocationAdapter.GetInvocationAdapter("json");
|
||||
|
||||
await SendRequest_IgnoreReceive(firstConnection, adapter, "ConnectionSendMethod", secondConnection.Connection.ConnectionId, "test");
|
||||
await firstClient.Invoke(nameof(MethodHub.ConnectionSendMethod), secondClient.Connection.ConnectionId, "test").OrTimeout();
|
||||
|
||||
// check that 'secondConnection' has received the group send
|
||||
var result = await ReadConnectionOutputAsync<InvocationDescriptor>(secondConnection).OrTimeout();
|
||||
var result = await secondClient.Read<InvocationDescriptor>().OrTimeout();
|
||||
Assert.Equal("Send", result.Method);
|
||||
Assert.Equal(1, result.Arguments.Length);
|
||||
Assert.Equal("test", result.Arguments[0]);
|
||||
|
||||
// kill the connections
|
||||
firstConnection.Connection.Dispose();
|
||||
secondConnection.Connection.Dispose();
|
||||
firstClient.Dispose();
|
||||
secondClient.Dispose();
|
||||
|
||||
await Task.WhenAll(firstEndPointTask, secondEndPointTask).OrTimeout();
|
||||
}
|
||||
|
|
@ -484,41 +440,6 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
return genericType.MakeGenericType(hubType);
|
||||
}
|
||||
|
||||
public async Task SendRequest(ConnectionWrapper connection, IInvocationAdapter writer, string method, params object[] args)
|
||||
{
|
||||
if (connection == null)
|
||||
{
|
||||
throw new ArgumentNullException();
|
||||
}
|
||||
|
||||
var stream = new MemoryStream();
|
||||
await writer.WriteMessageAsync(new InvocationDescriptor
|
||||
{
|
||||
Arguments = args,
|
||||
Method = method
|
||||
},
|
||||
stream);
|
||||
|
||||
var buffer = ReadableBuffer.Create(stream.ToArray()).Preserve();
|
||||
await connection.Application.Output.WriteAsync(new Message(buffer, Format.Binary, endOfMessage: true));
|
||||
}
|
||||
|
||||
public async Task SendRequest_IgnoreReceive(ConnectionWrapper connection, IInvocationAdapter writer, string method, params object[] args)
|
||||
{
|
||||
await SendRequest(connection, writer, method, args);
|
||||
|
||||
// Consume the result
|
||||
await connection.Application.Input.ReadAsync();
|
||||
}
|
||||
|
||||
private async Task<T> ReadConnectionOutputAsync<T>(ConnectionWrapper connection)
|
||||
{
|
||||
// TODO: other formats?
|
||||
var message = await connection.Application.Input.ReadAsync();
|
||||
var serializer = new JsonSerializer();
|
||||
return serializer.Deserialize<T>(new JsonTextReader(new StreamReader(new MemoryStream(message.Payload.Buffer.ToArray()))));
|
||||
}
|
||||
|
||||
private IServiceProvider CreateServiceProvider(Action<ServiceCollection> addServices = null)
|
||||
{
|
||||
var services = new ServiceCollection();
|
||||
|
|
@ -646,36 +567,5 @@ namespace Microsoft.AspNetCore.SignalR.Tests
|
|||
{
|
||||
public int DisposeCount = 0;
|
||||
}
|
||||
|
||||
public class ConnectionWrapper : IDisposable
|
||||
{
|
||||
private static int _id;
|
||||
|
||||
public Connection Connection { get; }
|
||||
|
||||
public IChannelConnection<Message> Application { get; }
|
||||
|
||||
public Task Connected => Connection.Metadata.Get<TaskCompletionSource<bool>>("ConnectedTask").Task;
|
||||
|
||||
public ConnectionWrapper(string format = "json")
|
||||
{
|
||||
var transportToApplication = Channel.CreateUnbounded<Message>();
|
||||
var applicationToTransport = Channel.CreateUnbounded<Message>();
|
||||
|
||||
Application = ChannelConnection.Create(input: applicationToTransport, output: transportToApplication);
|
||||
var transport = ChannelConnection.Create(input: transportToApplication, output: applicationToTransport);
|
||||
|
||||
Connection = new Connection(Guid.NewGuid().ToString(), transport);
|
||||
Connection.Metadata["formatType"] = format;
|
||||
Connection.User = new ClaimsPrincipal(new ClaimsIdentity(new[] { new Claim(ClaimTypes.Name, Interlocked.Increment(ref _id).ToString()) }));
|
||||
|
||||
Connection.Metadata["ConnectedTask"] = new TaskCompletionSource<bool>();
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
Connection.Dispose();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,120 @@
|
|||
// Copyright (c) .NET Foundation. All rights reserved.
|
||||
// Licensed under the Apache License, Version 2.0. See License.txt in the project root for license information.
|
||||
|
||||
using System;
|
||||
using System.IO;
|
||||
using System.IO.Pipelines;
|
||||
using System.Security.Claims;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using System.Threading.Tasks.Channels;
|
||||
using Microsoft.AspNetCore.Sockets;
|
||||
using Microsoft.AspNetCore.Sockets.Internal;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Microsoft.AspNetCore.SignalR.Tests
|
||||
{
|
||||
public class TestClient : IDisposable
|
||||
{
|
||||
private static int _id;
|
||||
private IInvocationAdapter _adapter;
|
||||
private CancellationTokenSource _cts;
|
||||
private TestBinder _binder;
|
||||
|
||||
public Connection Connection;
|
||||
public IChannelConnection<Message> Application { get; }
|
||||
public Task Connected => Connection.Metadata.Get<TaskCompletionSource<bool>>("ConnectedTask").Task;
|
||||
|
||||
public TestClient(IServiceProvider serviceProvider, string format = "json")
|
||||
{
|
||||
var transportToApplication = Channel.CreateUnbounded<Message>();
|
||||
var applicationToTransport = Channel.CreateUnbounded<Message>();
|
||||
|
||||
Application = ChannelConnection.Create(input: applicationToTransport, output: transportToApplication);
|
||||
var transport = ChannelConnection.Create(input: transportToApplication, output: applicationToTransport);
|
||||
|
||||
Connection = new Connection(Guid.NewGuid().ToString(), transport);
|
||||
Connection.Metadata["formatType"] = format;
|
||||
Connection.User = new ClaimsPrincipal(new ClaimsIdentity(new[] { new Claim(ClaimTypes.Name, Interlocked.Increment(ref _id).ToString()) }));
|
||||
Connection.Metadata["ConnectedTask"] = new TaskCompletionSource<bool>();
|
||||
|
||||
var invocationAdapter = serviceProvider.GetService<InvocationAdapterRegistry>();
|
||||
_adapter = invocationAdapter.GetInvocationAdapter(format);
|
||||
|
||||
_binder = new TestBinder();
|
||||
|
||||
_cts = new CancellationTokenSource();
|
||||
}
|
||||
|
||||
public async Task<T> Invoke<T>(string methodName, params object[] args) where T : InvocationMessage
|
||||
{
|
||||
await Invoke(methodName, args);
|
||||
|
||||
return await Read<T>();
|
||||
}
|
||||
|
||||
public async Task Invoke(string methodName, params object[] args)
|
||||
{
|
||||
var stream = new MemoryStream();
|
||||
await _adapter.WriteMessageAsync(new InvocationDescriptor
|
||||
{
|
||||
Arguments = args,
|
||||
Method = methodName
|
||||
},
|
||||
stream);
|
||||
|
||||
var buffer = ReadableBuffer.Create(stream.ToArray()).Preserve();
|
||||
await Application.Output.WriteAsync(new Message(buffer, Format.Binary, endOfMessage: true));
|
||||
}
|
||||
|
||||
public async Task<T> Read<T>() where T : InvocationMessage
|
||||
{
|
||||
while (await Application.Input.WaitToReadAsync(_cts.Token))
|
||||
{
|
||||
var value = await TryRead<T>();
|
||||
|
||||
if (value != null)
|
||||
{
|
||||
return value;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
public async Task<T> TryRead<T>() where T : InvocationMessage
|
||||
{
|
||||
Message message;
|
||||
if (Application.Input.TryRead(out message))
|
||||
{
|
||||
using (message)
|
||||
{
|
||||
var value = await _adapter.ReadMessageAsync(new MemoryStream(message.Payload.Buffer.ToArray()), _binder);
|
||||
return value as T;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
_cts.Cancel();
|
||||
Connection.Dispose();
|
||||
}
|
||||
|
||||
private class TestBinder : IInvocationBinder
|
||||
{
|
||||
public Type[] GetParameterTypes(string methodName)
|
||||
{
|
||||
// TODO: Possibly support actual client methods
|
||||
return new[] { typeof(object) };
|
||||
}
|
||||
|
||||
public Type GetReturnType(string invocationId)
|
||||
{
|
||||
return typeof(object);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue