integrate Channels-based WebSockets into SignalR (#28)
This commit is contained in:
parent
5e2b267d9f
commit
2431c5925c
|
|
@ -37,3 +37,5 @@ node_modules/
|
||||||
autobahnreports/
|
autobahnreports/
|
||||||
signalr-client.js
|
signalr-client.js
|
||||||
site.min.css
|
site.min.css
|
||||||
|
.idea/
|
||||||
|
.vscode/
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,3 @@
|
||||||
{
|
{
|
||||||
"projects": [ "src", "test" ],
|
"projects": [ "src", "test" ]
|
||||||
"sdk": {
|
|
||||||
"version": "1.0.0-preview2-003121"
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,6 +7,7 @@ using System.Threading.Tasks;
|
||||||
using Channels;
|
using Channels;
|
||||||
using Microsoft.AspNetCore.Http;
|
using Microsoft.AspNetCore.Http;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
|
using Microsoft.Extensions.Logging;
|
||||||
using Microsoft.Extensions.Primitives;
|
using Microsoft.Extensions.Primitives;
|
||||||
|
|
||||||
namespace Microsoft.AspNetCore.Sockets
|
namespace Microsoft.AspNetCore.Sockets
|
||||||
|
|
@ -15,11 +16,13 @@ namespace Microsoft.AspNetCore.Sockets
|
||||||
{
|
{
|
||||||
private readonly ConnectionManager _manager;
|
private readonly ConnectionManager _manager;
|
||||||
private readonly ChannelFactory _channelFactory;
|
private readonly ChannelFactory _channelFactory;
|
||||||
|
private readonly ILoggerFactory _loggerFactory;
|
||||||
|
|
||||||
public HttpConnectionDispatcher(ConnectionManager manager, ChannelFactory factory)
|
public HttpConnectionDispatcher(ConnectionManager manager, ChannelFactory factory, ILoggerFactory loggerFactory)
|
||||||
{
|
{
|
||||||
_manager = manager;
|
_manager = manager;
|
||||||
_channelFactory = factory;
|
_channelFactory = factory;
|
||||||
|
_loggerFactory = loggerFactory;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task ExecuteAsync<TEndPoint>(string path, HttpContext context) where TEndPoint : EndPoint
|
public async Task ExecuteAsync<TEndPoint>(string path, HttpContext context) where TEndPoint : EndPoint
|
||||||
|
|
@ -72,7 +75,7 @@ namespace Microsoft.AspNetCore.Sockets
|
||||||
var formatType = (string)context.Request.Query["formatType"];
|
var formatType = (string)context.Request.Query["formatType"];
|
||||||
state.Connection.Metadata["formatType"] = string.IsNullOrEmpty(formatType) ? "json" : formatType;
|
state.Connection.Metadata["formatType"] = string.IsNullOrEmpty(formatType) ? "json" : formatType;
|
||||||
|
|
||||||
var ws = new WebSockets(state.Connection, format);
|
var ws = new WebSockets(state.Connection, format, _loggerFactory);
|
||||||
|
|
||||||
await DoPersistentConnection(endpoint, ws, context, state.Connection);
|
await DoPersistentConnection(endpoint, ws, context, state.Connection);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -8,6 +8,7 @@ using Microsoft.AspNetCore.Routing;
|
||||||
using Microsoft.AspNetCore.Sockets;
|
using Microsoft.AspNetCore.Sockets;
|
||||||
using Microsoft.AspNetCore.Sockets.Routing;
|
using Microsoft.AspNetCore.Sockets.Routing;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
|
using Microsoft.Extensions.Logging;
|
||||||
|
|
||||||
namespace Microsoft.AspNetCore.Builder
|
namespace Microsoft.AspNetCore.Builder
|
||||||
{
|
{
|
||||||
|
|
@ -18,7 +19,8 @@ namespace Microsoft.AspNetCore.Builder
|
||||||
var manager = new ConnectionManager();
|
var manager = new ConnectionManager();
|
||||||
var factory = new ChannelFactory();
|
var factory = new ChannelFactory();
|
||||||
|
|
||||||
var dispatcher = new HttpConnectionDispatcher(manager, factory);
|
var loggerFactory = app.ApplicationServices.GetService<ILoggerFactory>();
|
||||||
|
var dispatcher = new HttpConnectionDispatcher(manager, factory, loggerFactory);
|
||||||
|
|
||||||
// Dispose the connection manager when application shutdown is triggered
|
// Dispose the connection manager when application shutdown is triggered
|
||||||
var lifetime = app.ApplicationServices.GetRequiredService<IApplicationLifetime>();
|
var lifetime = app.ApplicationServices.GetRequiredService<IApplicationLifetime>();
|
||||||
|
|
@ -28,8 +30,7 @@ namespace Microsoft.AspNetCore.Builder
|
||||||
|
|
||||||
callback(new SocketRouteBuilder(routes, dispatcher));
|
callback(new SocketRouteBuilder(routes, dispatcher));
|
||||||
|
|
||||||
// TODO: Use new low allocating websocket API
|
app.UseWebSocketConnections();
|
||||||
app.UseWebSockets();
|
|
||||||
app.UseRouter(routes.Build());
|
app.UseRouter(routes.Build());
|
||||||
return app;
|
return app;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,143 +2,155 @@
|
||||||
// Licensed under the Apache License, Version 2.0. See License.txt in the project root for license information.
|
// Licensed under the Apache License, Version 2.0. See License.txt in the project root for license information.
|
||||||
|
|
||||||
using System;
|
using System;
|
||||||
using System.Net.WebSockets;
|
|
||||||
using System.Threading;
|
|
||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
using Channels;
|
using Channels;
|
||||||
using Microsoft.AspNetCore.Http;
|
using Microsoft.AspNetCore.Http;
|
||||||
|
using Microsoft.AspNetCore.WebSockets.Internal;
|
||||||
|
using Microsoft.Extensions.Internal;
|
||||||
|
using Microsoft.Extensions.Logging;
|
||||||
|
using Microsoft.Extensions.Logging.Abstractions;
|
||||||
|
using Microsoft.Extensions.WebSockets.Internal;
|
||||||
|
|
||||||
namespace Microsoft.AspNetCore.Sockets
|
namespace Microsoft.AspNetCore.Sockets
|
||||||
{
|
{
|
||||||
public class WebSockets : IHttpTransport
|
public class WebSockets : IHttpTransport
|
||||||
{
|
{
|
||||||
private readonly HttpChannel _channel;
|
private static readonly TimeSpan _closeTimeout = TimeSpan.FromSeconds(5);
|
||||||
private readonly Connection _connection;
|
private static readonly WebSocketAcceptContext EmptyContext = new WebSocketAcceptContext();
|
||||||
private readonly WebSocketMessageType _messageType;
|
|
||||||
|
|
||||||
public WebSockets(Connection connection, Format format)
|
private readonly HttpChannel _channel;
|
||||||
|
private readonly WebSocketOpcode _opcode;
|
||||||
|
private readonly ILogger _logger;
|
||||||
|
|
||||||
|
public WebSockets(Connection connection, Format format, ILoggerFactory loggerFactory)
|
||||||
{
|
{
|
||||||
_connection = connection;
|
|
||||||
_channel = (HttpChannel)connection.Channel;
|
_channel = (HttpChannel)connection.Channel;
|
||||||
_messageType = format == Format.Binary ? WebSocketMessageType.Binary : WebSocketMessageType.Text;
|
_opcode = format == Format.Binary ? WebSocketOpcode.Binary : WebSocketOpcode.Text;
|
||||||
|
|
||||||
|
_logger = (ILogger)loggerFactory?.CreateLogger<WebSockets>() ?? NullLogger.Instance;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task ProcessRequestAsync(HttpContext context)
|
public async Task ProcessRequestAsync(HttpContext context)
|
||||||
{
|
{
|
||||||
if (!context.WebSockets.IsWebSocketRequest)
|
var feature = context.Features.Get<IHttpWebSocketConnectionFeature>();
|
||||||
|
if (feature == null || !feature.IsWebSocketRequest)
|
||||||
{
|
{
|
||||||
|
_logger.LogWarning("Unable to handle WebSocket request, there is no WebSocket feature available.");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
var ws = await context.WebSockets.AcceptWebSocketAsync();
|
using (var ws = await feature.AcceptWebSocketConnectionAsync(EmptyContext))
|
||||||
|
|
||||||
// REVIEW: Should we track this task? Leaving things like this alive usually causes memory leaks :)
|
|
||||||
// The reason we don't await this is because the channel is disposed after this loop returns
|
|
||||||
// and the sending loop is waiting for the channel to end before doing anything
|
|
||||||
// We could do a 2 stage shutdown but that could complicate the code...
|
|
||||||
var sending = StartSending(ws);
|
|
||||||
|
|
||||||
var outputBuffer = _channel.Input.Alloc();
|
|
||||||
|
|
||||||
while (!_channel.Input.Writing.IsCompleted)
|
|
||||||
{
|
{
|
||||||
// Make sure there's room to read (at least 2k)
|
_logger.LogInformation("Socket opened.");
|
||||||
outputBuffer.Ensure(2048);
|
|
||||||
|
|
||||||
ArraySegment<byte> segment;
|
// Begin sending and receiving. Receiving must be started first because ExecuteAsync enables SendAsync.
|
||||||
if (!outputBuffer.Memory.TryGetArray(out segment))
|
var receiving = ws.ExecuteAsync((frame, self) => ((WebSockets)self).HandleFrame(frame), this);
|
||||||
|
var sending = StartSending(ws);
|
||||||
|
|
||||||
|
// Wait for something to shut down.
|
||||||
|
var trigger = await Task.WhenAny(
|
||||||
|
receiving,
|
||||||
|
sending);
|
||||||
|
|
||||||
|
// What happened?
|
||||||
|
if (trigger == receiving)
|
||||||
{
|
{
|
||||||
// REVIEW: Do we care about native buffers here?
|
// Shutting down because we received a close frame from the client.
|
||||||
throw new InvalidOperationException("Managed buffers are required for Web Socket API");
|
// Complete the input writer so that the application knows there won't be any more input.
|
||||||
}
|
_logger.LogDebug("Client closed connection with status code '{0}' ({1}). Signaling end-of-input to application", receiving.Result.Status, receiving.Result.Description);
|
||||||
|
_channel.Input.CompleteWriter();
|
||||||
|
|
||||||
var result = await ws.ReceiveAsync(segment, CancellationToken.None);
|
// Wait for the application to finish sending.
|
||||||
|
_logger.LogDebug("Waiting for the application to finish sending data");
|
||||||
|
await sending;
|
||||||
|
|
||||||
if (result.MessageType != WebSocketMessageType.Close)
|
// Send the server's close frame
|
||||||
{
|
await ws.CloseAsync(WebSocketCloseStatus.NormalClosure);
|
||||||
outputBuffer.Advance(result.Count);
|
|
||||||
|
|
||||||
// Flush the written data to the channel
|
|
||||||
await outputBuffer.FlushAsync();
|
|
||||||
|
|
||||||
// Allocate a new buffer to further writing
|
|
||||||
outputBuffer = _channel.Input.Alloc();
|
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
break;
|
// The application finished sending. We're not going to keep the connection open,
|
||||||
|
// so close it and wait for the client to ack the close
|
||||||
|
_channel.Input.CompleteWriter();
|
||||||
|
_logger.LogDebug("Application finished sending. Sending close frame.");
|
||||||
|
await ws.CloseAsync(WebSocketCloseStatus.NormalClosure);
|
||||||
|
|
||||||
|
_logger.LogDebug("Waiting for the client to close the socket");
|
||||||
|
// TODO: Timeout.
|
||||||
|
await receiving;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
_logger.LogInformation("Socket closed.");
|
||||||
await ws.CloseOutputAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task StartSending(WebSocket ws)
|
private Task HandleFrame(WebSocketFrame frame)
|
||||||
{
|
{
|
||||||
while (true)
|
// Is this a frame we care about?
|
||||||
|
if (!frame.Opcode.IsMessage())
|
||||||
{
|
{
|
||||||
var result = await _channel.Output.ReadAsync();
|
return TaskCache.CompletedTask;
|
||||||
var buffer = result.Buffer;
|
|
||||||
|
|
||||||
try
|
|
||||||
{
|
|
||||||
if (buffer.IsEmpty && result.IsCompleted)
|
|
||||||
{
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
foreach (var memory in buffer)
|
|
||||||
{
|
|
||||||
ArraySegment<byte> data;
|
|
||||||
if (memory.TryGetArray(out data))
|
|
||||||
{
|
|
||||||
if (IsClosedOrClosedSent(ws))
|
|
||||||
{
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
await ws.SendAsync(data, _messageType, endOfMessage: true, cancellationToken: CancellationToken.None);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
catch (Exception)
|
|
||||||
{
|
|
||||||
// Error writing, probably closed
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
finally
|
|
||||||
{
|
|
||||||
_channel.Output.Advance(buffer.End);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// REVIEW: Should this ever happen?
|
LogFrame("Receiving", frame);
|
||||||
if (!IsClosedOrClosedSent(ws))
|
|
||||||
{
|
// Allocate space from the input channel
|
||||||
// Close the output
|
var outputBuffer = _channel.Input.Alloc();
|
||||||
await ws.CloseOutputAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None);
|
|
||||||
}
|
// Append this buffer to the input channel
|
||||||
|
_logger.LogDebug($"Appending {frame.Payload.Length} bytes to Connection channel");
|
||||||
|
outputBuffer.Append(frame.Payload);
|
||||||
|
|
||||||
|
return outputBuffer.FlushAsync();
|
||||||
}
|
}
|
||||||
|
|
||||||
private static bool IsClosedOrClosedSent(WebSocket webSocket)
|
private void LogFrame(string action, WebSocketFrame frame)
|
||||||
{
|
{
|
||||||
var webSocketState = GetWebSocketState(webSocket);
|
if (_logger.IsEnabled(LogLevel.Debug))
|
||||||
|
{
|
||||||
return webSocketState == WebSocketState.Closed ||
|
_logger.LogDebug(
|
||||||
webSocketState == WebSocketState.CloseSent ||
|
$"{action} frame: Opcode={frame.Opcode}, Fin={frame.EndOfMessage}, Payload={frame.Payload.Length} bytes");
|
||||||
webSocketState == WebSocketState.Aborted;
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private static WebSocketState GetWebSocketState(WebSocket webSocket)
|
private async Task StartSending(IWebSocketConnection ws)
|
||||||
{
|
{
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
return webSocket.State;
|
while (true)
|
||||||
|
{
|
||||||
|
var result = await _channel.Output.ReadAsync();
|
||||||
|
var buffer = result.Buffer;
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
if (buffer.IsEmpty && result.IsCompleted)
|
||||||
|
{
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Send the buffer in a frame
|
||||||
|
var frame = new WebSocketFrame(
|
||||||
|
endOfMessage: true,
|
||||||
|
opcode: _opcode,
|
||||||
|
payload: buffer);
|
||||||
|
LogFrame("Sending", frame);
|
||||||
|
await ws.SendAsync(frame);
|
||||||
|
}
|
||||||
|
catch (Exception ex)
|
||||||
|
{
|
||||||
|
_logger.LogError("Error writing frame to output: {0}", ex);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
_channel.Output.Advance(buffer.End);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
catch (ObjectDisposedException)
|
finally
|
||||||
{
|
{
|
||||||
return WebSocketState.Closed;
|
// No longer reading from the channel
|
||||||
|
_channel.Output.CompleteReader();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -25,7 +25,11 @@
|
||||||
"Channels": "0.2.0-beta-*",
|
"Channels": "0.2.0-beta-*",
|
||||||
"Microsoft.AspNetCore.Hosting.Abstractions": "1.2.0-*",
|
"Microsoft.AspNetCore.Hosting.Abstractions": "1.2.0-*",
|
||||||
"Microsoft.AspNetCore.Routing": "1.2.0-*",
|
"Microsoft.AspNetCore.Routing": "1.2.0-*",
|
||||||
"Microsoft.AspNetCore.WebSockets": "1.1.0-*",
|
"Microsoft.AspNetCore.WebSockets.Internal": "0.1.0-*",
|
||||||
|
"Microsoft.Extensions.TaskCache.Sources": {
|
||||||
|
"version": "1.2.0-*",
|
||||||
|
"type": "build"
|
||||||
|
},
|
||||||
"NETStandard.Library": "1.6.1-*"
|
"NETStandard.Library": "1.6.1-*"
|
||||||
},
|
},
|
||||||
"frameworks": {
|
"frameworks": {
|
||||||
|
|
|
||||||
|
|
@ -34,7 +34,6 @@ namespace Microsoft.AspNetCore.Builder
|
||||||
{
|
{
|
||||||
throw new ArgumentNullException(nameof(options));
|
throw new ArgumentNullException(nameof(options));
|
||||||
}
|
}
|
||||||
self.UseWebSocketConnections(channelFactory, options);
|
|
||||||
self.UseMiddleware<WebSocketConnectionMiddleware>(channelFactory, options);
|
self.UseMiddleware<WebSocketConnectionMiddleware>(channelFactory, options);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,7 @@ namespace Microsoft.AspNetCore.Sockets.Tests
|
||||||
var manager = new ConnectionManager();
|
var manager = new ConnectionManager();
|
||||||
using (var factory = new ChannelFactory())
|
using (var factory = new ChannelFactory())
|
||||||
{
|
{
|
||||||
var dispatcher = new HttpConnectionDispatcher(manager, factory);
|
var dispatcher = new HttpConnectionDispatcher(manager, factory, loggerFactory: null);
|
||||||
var context = new DefaultHttpContext();
|
var context = new DefaultHttpContext();
|
||||||
var ms = new MemoryStream();
|
var ms = new MemoryStream();
|
||||||
context.Request.Path = "/getid";
|
context.Request.Path = "/getid";
|
||||||
|
|
@ -46,7 +46,7 @@ namespace Microsoft.AspNetCore.Sockets.Tests
|
||||||
|
|
||||||
using (var factory = new ChannelFactory())
|
using (var factory = new ChannelFactory())
|
||||||
{
|
{
|
||||||
var dispatcher = new HttpConnectionDispatcher(manager, factory);
|
var dispatcher = new HttpConnectionDispatcher(manager, factory, loggerFactory: null);
|
||||||
var context = new DefaultHttpContext();
|
var context = new DefaultHttpContext();
|
||||||
context.Request.Path = "/send";
|
context.Request.Path = "/send";
|
||||||
var values = new Dictionary<string, StringValues>();
|
var values = new Dictionary<string, StringValues>();
|
||||||
|
|
@ -66,7 +66,7 @@ namespace Microsoft.AspNetCore.Sockets.Tests
|
||||||
var manager = new ConnectionManager();
|
var manager = new ConnectionManager();
|
||||||
using (var factory = new ChannelFactory())
|
using (var factory = new ChannelFactory())
|
||||||
{
|
{
|
||||||
var dispatcher = new HttpConnectionDispatcher(manager, factory);
|
var dispatcher = new HttpConnectionDispatcher(manager, factory, loggerFactory: null);
|
||||||
var context = new DefaultHttpContext();
|
var context = new DefaultHttpContext();
|
||||||
context.Request.Path = "/send";
|
context.Request.Path = "/send";
|
||||||
var values = new Dictionary<string, StringValues>();
|
var values = new Dictionary<string, StringValues>();
|
||||||
|
|
@ -86,7 +86,7 @@ namespace Microsoft.AspNetCore.Sockets.Tests
|
||||||
var manager = new ConnectionManager();
|
var manager = new ConnectionManager();
|
||||||
using (var factory = new ChannelFactory())
|
using (var factory = new ChannelFactory())
|
||||||
{
|
{
|
||||||
var dispatcher = new HttpConnectionDispatcher(manager, factory);
|
var dispatcher = new HttpConnectionDispatcher(manager, factory, loggerFactory: null);
|
||||||
var context = new DefaultHttpContext();
|
var context = new DefaultHttpContext();
|
||||||
context.Request.Path = "/send";
|
context.Request.Path = "/send";
|
||||||
await Assert.ThrowsAsync<InvalidOperationException>(async () =>
|
await Assert.ThrowsAsync<InvalidOperationException>(async () =>
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
{
|
{
|
||||||
"buildOptions": {
|
"buildOptions": {
|
||||||
"warningsAsErrors": true
|
"warningsAsErrors": true
|
||||||
},
|
},
|
||||||
|
|
@ -7,6 +7,8 @@
|
||||||
"dotnet-test-xunit": "2.2.0-*",
|
"dotnet-test-xunit": "2.2.0-*",
|
||||||
"Microsoft.AspNetCore.Http": "1.2.0-*",
|
"Microsoft.AspNetCore.Http": "1.2.0-*",
|
||||||
"Microsoft.AspNetCore.Sockets": "0.1.0-*",
|
"Microsoft.AspNetCore.Sockets": "0.1.0-*",
|
||||||
|
"Microsoft.AspNetCore.Hosting": "1.1.0-*",
|
||||||
|
"Microsoft.AspNetCore.Server.Kestrel": "1.1.0-*",
|
||||||
"xunit": "2.2.0-*"
|
"xunit": "2.2.0-*"
|
||||||
},
|
},
|
||||||
"frameworks": {
|
"frameworks": {
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue