Use TryRead and TryWrite (#113)
* Use TryRead and TryWrite - Use TryWrite to avoid errors on channel close for /send requests - Use TryRead until it returns false for all transports but long polling
This commit is contained in:
parent
5d374b7dbe
commit
8dc68cb798
|
|
@ -145,58 +145,61 @@ namespace Microsoft.AspNetCore.SignalR
|
||||||
|
|
||||||
while (await connection.Transport.Input.WaitToReadAsync())
|
while (await connection.Transport.Input.WaitToReadAsync())
|
||||||
{
|
{
|
||||||
Message message;
|
Message incomingMessage;
|
||||||
if (!connection.Transport.Input.TryRead(out message))
|
while (connection.Transport.Input.TryRead(out incomingMessage))
|
||||||
{
|
{
|
||||||
continue;
|
InvocationDescriptor invocationDescriptor;
|
||||||
}
|
using (incomingMessage)
|
||||||
|
|
||||||
InvocationDescriptor invocationDescriptor;
|
|
||||||
using (message)
|
|
||||||
{
|
|
||||||
var inputStream = new MemoryStream(message.Payload.Buffer.ToArray());
|
|
||||||
|
|
||||||
// TODO: Handle receiving InvocationResultDescriptor
|
|
||||||
invocationDescriptor = await invocationAdapter.ReadMessageAsync(inputStream, this) as InvocationDescriptor;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Is there a better way of detecting that a connection was closed?
|
|
||||||
if (invocationDescriptor == null)
|
|
||||||
{
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (_logger.IsEnabled(LogLevel.Debug))
|
|
||||||
{
|
|
||||||
_logger.LogDebug("Received hub invocation: {invocation}", invocationDescriptor);
|
|
||||||
}
|
|
||||||
|
|
||||||
InvocationResultDescriptor result;
|
|
||||||
Func<Connection, InvocationDescriptor, Task<InvocationResultDescriptor>> callback;
|
|
||||||
if (_callbacks.TryGetValue(invocationDescriptor.Method, out callback))
|
|
||||||
{
|
|
||||||
result = await callback(connection, invocationDescriptor);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
// If there's no method then return a failed response for this request
|
|
||||||
result = new InvocationResultDescriptor
|
|
||||||
{
|
{
|
||||||
Id = invocationDescriptor.Id,
|
var inputStream = new MemoryStream(incomingMessage.Payload.Buffer.ToArray());
|
||||||
Error = $"Unknown hub method '{invocationDescriptor.Method}'"
|
|
||||||
};
|
|
||||||
|
|
||||||
_logger.LogError("Unknown hub method '{method}'", invocationDescriptor.Method);
|
// TODO: Handle receiving InvocationResultDescriptor
|
||||||
}
|
invocationDescriptor = await invocationAdapter.ReadMessageAsync(inputStream, this) as InvocationDescriptor;
|
||||||
|
}
|
||||||
|
|
||||||
// TODO: Pool memory
|
// Is there a better way of detecting that a connection was closed?
|
||||||
var outStream = new MemoryStream();
|
if (invocationDescriptor == null)
|
||||||
await invocationAdapter.WriteMessageAsync(result, outStream);
|
{
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
var buffer = ReadableBuffer.Create(outStream.ToArray()).Preserve();
|
if (_logger.IsEnabled(LogLevel.Debug))
|
||||||
if (await connection.Transport.Output.WaitToWriteAsync())
|
{
|
||||||
{
|
_logger.LogDebug("Received hub invocation: {invocation}", invocationDescriptor);
|
||||||
connection.Transport.Output.TryWrite(new Message(buffer, connection.Metadata.Format, endOfMessage: true));
|
}
|
||||||
|
|
||||||
|
InvocationResultDescriptor result;
|
||||||
|
Func<Connection, InvocationDescriptor, Task<InvocationResultDescriptor>> callback;
|
||||||
|
if (_callbacks.TryGetValue(invocationDescriptor.Method, out callback))
|
||||||
|
{
|
||||||
|
result = await callback(connection, invocationDescriptor);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
// If there's no method then return a failed response for this request
|
||||||
|
result = new InvocationResultDescriptor
|
||||||
|
{
|
||||||
|
Id = invocationDescriptor.Id,
|
||||||
|
Error = $"Unknown hub method '{invocationDescriptor.Method}'"
|
||||||
|
};
|
||||||
|
|
||||||
|
_logger.LogError("Unknown hub method '{method}'", invocationDescriptor.Method);
|
||||||
|
}
|
||||||
|
|
||||||
|
// TODO: Pool memory
|
||||||
|
var outStream = new MemoryStream();
|
||||||
|
await invocationAdapter.WriteMessageAsync(result, outStream);
|
||||||
|
|
||||||
|
var buffer = ReadableBuffer.Create(outStream.ToArray()).Preserve();
|
||||||
|
var outMessage = new Message(buffer, connection.Metadata.Format, endOfMessage: true);
|
||||||
|
|
||||||
|
while (await connection.Transport.Output.WaitToWriteAsync())
|
||||||
|
{
|
||||||
|
if (connection.Transport.Output.TryWrite(outMessage))
|
||||||
|
{
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -225,8 +225,14 @@ namespace Microsoft.AspNetCore.Sockets
|
||||||
format,
|
format,
|
||||||
endOfMessage: true);
|
endOfMessage: true);
|
||||||
|
|
||||||
await state.Application.Output.WriteAsync(message);
|
// REVIEW: Do we want to return a specific status code here if the connection has ended?
|
||||||
|
while (await state.Application.Output.WaitToWriteAsync())
|
||||||
|
{
|
||||||
|
if (state.Application.Output.TryWrite(message))
|
||||||
|
{
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -24,20 +24,17 @@ namespace Microsoft.AspNetCore.Sockets.Transports
|
||||||
|
|
||||||
public async Task ProcessRequestAsync(HttpContext context)
|
public async Task ProcessRequestAsync(HttpContext context)
|
||||||
{
|
{
|
||||||
if (_application.Completion.IsCompleted)
|
|
||||||
{
|
|
||||||
// Client should stop if it receives a 204
|
|
||||||
_logger.LogInformation("Terminating Long Polling connection by sending 204 response.");
|
|
||||||
context.Response.StatusCode = 204;
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
// TODO: We need the ability to yield the connection without completing the channel.
|
// TODO: We need the ability to yield the connection without completing the channel.
|
||||||
// This is to force ReadAsync to yield without data to end to poll but not the entire connection.
|
// This is to force ReadAsync to yield without data to end to poll but not the entire connection.
|
||||||
// This is for cases when the client reconnects see issue #27
|
// This is for cases when the client reconnects see issue #27
|
||||||
await _application.WaitToReadAsync(context.RequestAborted);
|
if (!await _application.WaitToReadAsync(context.RequestAborted))
|
||||||
|
{
|
||||||
|
_logger.LogInformation("Terminating Long Polling connection by sending 204 response.");
|
||||||
|
context.Response.StatusCode = 204;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
Message message;
|
Message message;
|
||||||
if (_application.TryRead(out message))
|
if (_application.TryRead(out message))
|
||||||
|
|
|
||||||
|
|
@ -33,7 +33,7 @@ namespace Microsoft.AspNetCore.Sockets.Transports
|
||||||
while (await _application.WaitToReadAsync(context.RequestAborted))
|
while (await _application.WaitToReadAsync(context.RequestAborted))
|
||||||
{
|
{
|
||||||
Message message;
|
Message message;
|
||||||
if (_application.TryRead(out message))
|
while (_application.TryRead(out message))
|
||||||
{
|
{
|
||||||
using (message)
|
using (message)
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -156,7 +156,7 @@ namespace Microsoft.AspNetCore.Sockets.Transports
|
||||||
{
|
{
|
||||||
// Get a frame from the application
|
// Get a frame from the application
|
||||||
Message message;
|
Message message;
|
||||||
if (_application.Input.TryRead(out message))
|
while (_application.Input.TryRead(out message))
|
||||||
{
|
{
|
||||||
using (message)
|
using (message)
|
||||||
{
|
{
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue