Clean up logging (#1308)

This commit is contained in:
BrennanConroy 2018-01-22 09:37:53 -08:00 committed by GitHub
parent 0579f40a7d
commit a449345436
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
27 changed files with 616 additions and 776 deletions

View File

@ -174,6 +174,12 @@ export class HttpConnection implements IConnection {
this.transport = null; this.transport = null;
} }
if (error) {
this.logger.log(LogLevel.Error, `Connection disconnected with error '${error}'.`);
} else {
this.logger.log(LogLevel.Information, "Connection disconnected.");
}
this.connectionState = ConnectionState.Disconnected; this.connectionState = ConnectionState.Disconnected;
if (raiseClosed && this.onclose) { if (raiseClosed && this.onclose) {

View File

@ -132,6 +132,8 @@ export class HubConnection {
this.callbacks.clear(); this.callbacks.clear();
this.closedCallbacks.forEach(c => c.apply(this, [error])); this.closedCallbacks.forEach(c => c.apply(this, [error]));
this.cleanupTimeout();
} }
async start(): Promise<void> { async start(): Promise<void> {
@ -158,9 +160,7 @@ export class HubConnection {
} }
stop(): Promise<void> { stop(): Promise<void> {
if (this.timeoutHandle) { this.cleanupTimeout();
clearTimeout(this.timeoutHandle);
}
return this.connection.stop(); return this.connection.stop();
} }
@ -285,6 +285,12 @@ export class HubConnection {
} }
} }
private cleanupTimeout(): void {
if (this.timeoutHandle) {
clearTimeout(this.timeoutHandle);
}
}
private createInvocation(methodName: string, args: any[], nonblocking: boolean): InvocationMessage { private createInvocation(methodName: string, args: any[], nonblocking: boolean): InvocationMessage {
if (nonblocking) { if (nonblocking) {
return { return {

View File

@ -104,7 +104,7 @@
connectButton.disabled = true; connectButton.disabled = true;
disconnectButton.disabled = false; disconnectButton.disabled = false;
console.log('http://' + document.location.host + '/' + hubRoute); console.log('http://' + document.location.host + '/' + hubRoute);
connection = new signalR.HubConnection(hubRoute, { transport: transportType, logging: logger }); connection = new signalR.HubConnection(hubRoute, { transport: transportType, logger: logger });
connection.on('Send', function (msg) { connection.on('Send', function (msg) {
addLine('message-list', msg); addLine('message-list', msg);
}); });

View File

@ -1,4 +1,4 @@
<!DOCTYPE html> <!DOCTYPE html>
<html> <html>
<head> <head>
@ -24,7 +24,7 @@
document.getElementById('transportName').innerHTML = signalR.TransportType[transportType]; document.getElementById('transportName').innerHTML = signalR.TransportType[transportType];
let url = 'http://' + document.location.host + '/chat'; let url = 'http://' + document.location.host + '/chat';
let connection = new signalR.HttpConnection(url, { transport: transportType, logging: new signalR.ConsoleLogger(signalR.LogLevel.Trace) }); let connection = new signalR.HttpConnection(url, { transport: transportType, logger: new signalR.ConsoleLogger(signalR.LogLevel.Trace) });
connection.onreceive = function (data) { connection.onreceive = function (data) {
let child = document.createElement('li'); let child = document.createElement('li');

View File

@ -1,4 +1,4 @@
<!DOCTYPE html> <!DOCTYPE html>
<html> <html>
<head> <head>
@ -54,7 +54,7 @@
}); });
click('connectButton', function () { click('connectButton', function () {
connection = new signalR.HubConnection('/streaming', { transport: transportType, logging: logger }); connection = new signalR.HubConnection('/streaming', { transport: transportType, logger: logger });
connection.onclose(function () { connection.onclose(function () {
channelButton.disabled = true; channelButton.disabled = true;

View File

@ -313,7 +313,7 @@ namespace Microsoft.AspNetCore.SignalR.Client
try try
{ {
_logger.PreparingNonBlockingInvocation(invocationMessage.InvocationId, methodName, args.Length); _logger.PreparingNonBlockingInvocation(methodName, args.Length);
var payload = _protocolReaderWriter.WriteMessage(invocationMessage); var payload = _protocolReaderWriter.WriteMessage(invocationMessage);
_logger.SendInvocation(invocationMessage.InvocationId); _logger.SendInvocation(invocationMessage.InvocationId);

View File

@ -10,8 +10,8 @@ namespace Microsoft.AspNetCore.SignalR.Client.Internal
internal static class SignalRClientLoggerExtensions internal static class SignalRClientLoggerExtensions
{ {
// Category: HubConnection // Category: HubConnection
private static readonly Action<ILogger, string, string, int, Exception> _preparingNonBlockingInvocation = private static readonly Action<ILogger, string, int, Exception> _preparingNonBlockingInvocation =
LoggerMessage.Define<string, string, int>(LogLevel.Trace, new EventId(1, nameof(PreparingNonBlockingInvocation)), "Preparing non-blocking invocation '{invocationId}' of '{target}', with {argumentCount} argument(s)."); LoggerMessage.Define<string, int>(LogLevel.Trace, new EventId(1, nameof(PreparingNonBlockingInvocation)), "Preparing non-blocking invocation of '{target}', with {argumentCount} argument(s).");
private static readonly Action<ILogger, string, string, string, int, Exception> _preparingBlockingInvocation = private static readonly Action<ILogger, string, string, string, int, Exception> _preparingBlockingInvocation =
LoggerMessage.Define<string, string, string, int>(LogLevel.Trace, new EventId(2, nameof(PreparingBlockingInvocation)), "Preparing blocking invocation '{invocationId}' of '{target}', with return type '{returnType}' and {argumentCount} argument(s)."); LoggerMessage.Define<string, string, string, int>(LogLevel.Trace, new EventId(2, nameof(PreparingBlockingInvocation)), "Preparing blocking invocation '{invocationId}' of '{target}', with return type '{returnType}' and {argumentCount} argument(s).");
@ -118,9 +118,9 @@ namespace Microsoft.AspNetCore.SignalR.Client.Internal
private static readonly Action<ILogger, string, Exception> _errorInvokingClientSideMethod = private static readonly Action<ILogger, string, Exception> _errorInvokingClientSideMethod =
LoggerMessage.Define<string>(LogLevel.Error, new EventId(6, nameof(ErrorInvokingClientSideMethod)), "Invoking client side method '{methodName}' failed."); LoggerMessage.Define<string>(LogLevel.Error, new EventId(6, nameof(ErrorInvokingClientSideMethod)), "Invoking client side method '{methodName}' failed.");
public static void PreparingNonBlockingInvocation(this ILogger logger, string invocationId, string target, int count) public static void PreparingNonBlockingInvocation(this ILogger logger, string target, int count)
{ {
_preparingNonBlockingInvocation(logger, invocationId, target, count, null); _preparingNonBlockingInvocation(logger, target, count, null);
} }
public static void PreparingBlockingInvocation(this ILogger logger, string invocationId, string target, string returnType, int count) public static void PreparingBlockingInvocation(this ILogger logger, string invocationId, string target, string returnType, int count)

View File

@ -306,13 +306,13 @@ namespace Microsoft.AspNetCore.SignalR
if (isStreamedInvocation) if (isStreamedInvocation)
{ {
var enumerator = GetStreamingEnumerator(connection, hubMethodInvocationMessage.InvocationId, methodExecutor, result, methodExecutor.MethodReturnType); var enumerator = GetStreamingEnumerator(connection, hubMethodInvocationMessage.InvocationId, methodExecutor, result, methodExecutor.MethodReturnType);
_logger.StreamingResult(hubMethodInvocationMessage.InvocationId, methodExecutor.MethodReturnType.FullName); _logger.StreamingResult(hubMethodInvocationMessage.InvocationId, methodExecutor);
await StreamResultsAsync(hubMethodInvocationMessage.InvocationId, connection, enumerator); await StreamResultsAsync(hubMethodInvocationMessage.InvocationId, connection, enumerator);
} }
// Non-empty/null InvocationId ==> Blocking invocation that needs a response // Non-empty/null InvocationId ==> Blocking invocation that needs a response
else if (!string.IsNullOrEmpty(hubMethodInvocationMessage.InvocationId)) else if (!string.IsNullOrEmpty(hubMethodInvocationMessage.InvocationId))
{ {
_logger.SendingResult(hubMethodInvocationMessage.InvocationId, methodExecutor.MethodReturnType.FullName); _logger.SendingResult(hubMethodInvocationMessage.InvocationId, methodExecutor);
await SendMessageAsync(connection, CompletionMessage.WithResult(hubMethodInvocationMessage.InvocationId, result)); await SendMessageAsync(connection, CompletionMessage.WithResult(hubMethodInvocationMessage.InvocationId, result));
} }
} }
@ -506,6 +506,7 @@ namespace Microsoft.AspNetCore.SignalR
{ {
var hubType = typeof(THub); var hubType = typeof(THub);
var hubTypeInfo = hubType.GetTypeInfo(); var hubTypeInfo = hubType.GetTypeInfo();
var hubName = hubType.Name;
foreach (var methodInfo in HubReflectionHelper.GetHubMethods(hubType)) foreach (var methodInfo in HubReflectionHelper.GetHubMethods(hubType))
{ {
@ -522,7 +523,7 @@ namespace Microsoft.AspNetCore.SignalR
var authorizeAttributes = methodInfo.GetCustomAttributes<AuthorizeAttribute>(inherit: true); var authorizeAttributes = methodInfo.GetCustomAttributes<AuthorizeAttribute>(inherit: true);
_methods[methodName] = new HubMethodDescriptor(executor, authorizeAttributes); _methods[methodName] = new HubMethodDescriptor(executor, authorizeAttributes);
_logger.HubMethodBound(methodName); _logger.HubMethodBound(hubName, methodName);
} }
} }

View File

@ -3,6 +3,7 @@
using System; using System;
using Microsoft.AspNetCore.SignalR.Internal.Protocol; using Microsoft.AspNetCore.SignalR.Internal.Protocol;
using Microsoft.Extensions.Internal;
using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging;
namespace Microsoft.AspNetCore.SignalR.Internal namespace Microsoft.AspNetCore.SignalR.Internal
@ -32,16 +33,16 @@ namespace Microsoft.AspNetCore.SignalR.Internal
LoggerMessage.Define<string>(LogLevel.Debug, new EventId(7, nameof(HubMethodNotAuthorized)), "Failed to invoke '{hubMethod}' because user is unauthorized."); LoggerMessage.Define<string>(LogLevel.Debug, new EventId(7, nameof(HubMethodNotAuthorized)), "Failed to invoke '{hubMethod}' because user is unauthorized.");
private static readonly Action<ILogger, string, string, Exception> _streamingResult = private static readonly Action<ILogger, string, string, Exception> _streamingResult =
LoggerMessage.Define<string, string>(LogLevel.Trace, new EventId(8, nameof(StreamingResult)), "{invocationId}: Streaming result of type '{resultType}'."); LoggerMessage.Define<string, string>(LogLevel.Trace, new EventId(8, nameof(StreamingResult)), "InvocationId {invocationId}: Streaming result of type '{resultType}'.");
private static readonly Action<ILogger, string, string, Exception> _sendingResult = private static readonly Action<ILogger, string, string, Exception> _sendingResult =
LoggerMessage.Define<string, string>(LogLevel.Trace, new EventId(9, nameof(SendingResult)), "{invocationId}: Sending result of type '{resultType}'."); LoggerMessage.Define<string, string>(LogLevel.Trace, new EventId(9, nameof(SendingResult)), "InvocationId {invocationId}: Sending result of type '{resultType}'.");
private static readonly Action<ILogger, string, Exception> _failedInvokingHubMethod = private static readonly Action<ILogger, string, Exception> _failedInvokingHubMethod =
LoggerMessage.Define<string>(LogLevel.Error, new EventId(10, nameof(FailedInvokingHubMethod)), "Failed to invoke hub method '{hubMethod}'."); LoggerMessage.Define<string>(LogLevel.Error, new EventId(10, nameof(FailedInvokingHubMethod)), "Failed to invoke hub method '{hubMethod}'.");
private static readonly Action<ILogger, string, Exception> _hubMethodBound = private static readonly Action<ILogger, string, string, Exception> _hubMethodBound =
LoggerMessage.Define<string>(LogLevel.Trace, new EventId(11, nameof(HubMethodBound)), "Hub method '{hubMethod}' is bound."); LoggerMessage.Define<string, string>(LogLevel.Trace, new EventId(11, nameof(HubMethodBound)), "'{hubName}' hub method '{hubMethod}' is bound.");
private static readonly Action<ILogger, string, Exception> _cancelStream = private static readonly Action<ILogger, string, Exception> _cancelStream =
LoggerMessage.Define<string>(LogLevel.Debug, new EventId(12, nameof(CancelStream)), "Canceling stream for invocation {invocationId}."); LoggerMessage.Define<string>(LogLevel.Debug, new EventId(12, nameof(CancelStream)), "Canceling stream for invocation {invocationId}.");
@ -129,14 +130,16 @@ namespace Microsoft.AspNetCore.SignalR.Internal
_hubMethodNotAuthorized(logger, hubMethod, null); _hubMethodNotAuthorized(logger, hubMethod, null);
} }
public static void StreamingResult(this ILogger logger, string invocationId, string resultType) public static void StreamingResult(this ILogger logger, string invocationId, ObjectMethodExecutor objectMethodExecutor)
{ {
_streamingResult(logger, invocationId, resultType, null); var resultType = objectMethodExecutor.AsyncResultType == null ? objectMethodExecutor.MethodReturnType : objectMethodExecutor.AsyncResultType;
_streamingResult(logger, invocationId, resultType.FullName, null);
} }
public static void SendingResult(this ILogger logger, string invocationId, string resultType) public static void SendingResult(this ILogger logger, string invocationId, ObjectMethodExecutor objectMethodExecutor)
{ {
_sendingResult(logger, invocationId, resultType, null); var resultType = objectMethodExecutor.AsyncResultType == null ? objectMethodExecutor.MethodReturnType : objectMethodExecutor.AsyncResultType;
_sendingResult(logger, invocationId, resultType.FullName, null);
} }
public static void FailedInvokingHubMethod(this ILogger logger, string hubMethod, Exception exception) public static void FailedInvokingHubMethod(this ILogger logger, string hubMethod, Exception exception)
@ -144,9 +147,9 @@ namespace Microsoft.AspNetCore.SignalR.Internal
_failedInvokingHubMethod(logger, hubMethod, exception); _failedInvokingHubMethod(logger, hubMethod, exception);
} }
public static void HubMethodBound(this ILogger logger, string hubMethod) public static void HubMethodBound(this ILogger logger, string hubName, string hubMethod)
{ {
_hubMethodBound(logger, hubMethod, null); _hubMethodBound(logger, hubName, hubMethod, null);
} }
public static void CancelStream(this ILogger logger, string invocationId) public static void CancelStream(this ILogger logger, string invocationId)

View File

@ -47,6 +47,8 @@ namespace Microsoft.AspNetCore.Sockets.Client
private ChannelWriter<SendMessage> Output => _transportChannel.Output; private ChannelWriter<SendMessage> Output => _transportChannel.Output;
private readonly List<ReceiveCallback> _callbacks = new List<ReceiveCallback>(); private readonly List<ReceiveCallback> _callbacks = new List<ReceiveCallback>();
private readonly TransportType _requestedTransportType = TransportType.All; private readonly TransportType _requestedTransportType = TransportType.All;
private readonly ConnectionLogScope _logScope;
private readonly IDisposable _scopeDisposable;
public Uri Url { get; } public Uri Url { get; }
@ -89,6 +91,8 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
_transportFactory = new DefaultTransportFactory(transportType, _loggerFactory, _httpClient, httpOptions); _transportFactory = new DefaultTransportFactory(transportType, _loggerFactory, _httpClient, httpOptions);
_logScope = new ConnectionLogScope();
_scopeDisposable = _logger.BeginScope(_logScope);
} }
public HttpConnection(Uri url, ITransportFactory transportFactory, ILoggerFactory loggerFactory, HttpOptions httpOptions) public HttpConnection(Uri url, ITransportFactory transportFactory, ILoggerFactory loggerFactory, HttpOptions httpOptions)
@ -100,6 +104,8 @@ namespace Microsoft.AspNetCore.Sockets.Client
_httpClient = _httpOptions?.HttpMessageHandler == null ? new HttpClient() : new HttpClient(_httpOptions?.HttpMessageHandler); _httpClient = _httpOptions?.HttpMessageHandler == null ? new HttpClient() : new HttpClient(_httpOptions?.HttpMessageHandler);
_httpClient.Timeout = HttpClientTimeout; _httpClient.Timeout = HttpClientTimeout;
_transportFactory = transportFactory ?? throw new ArgumentNullException(nameof(transportFactory)); _transportFactory = transportFactory ?? throw new ArgumentNullException(nameof(transportFactory));
_logScope = new ConnectionLogScope();
_scopeDisposable = _logger.BeginScope(_logScope);
} }
public async Task StartAsync() => await StartAsyncCore().ForceAsync(); public async Task StartAsync() => await StartAsyncCore().ForceAsync();
@ -151,11 +157,12 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
var negotiationResponse = await Negotiate(Url, _httpClient, _logger); var negotiationResponse = await Negotiate(Url, _httpClient, _logger);
_connectionId = negotiationResponse.ConnectionId; _connectionId = negotiationResponse.ConnectionId;
_logScope.ConnectionId = _connectionId;
// Connection is being disposed while start was in progress // Connection is being disposed while start was in progress
if (_connectionState == ConnectionState.Disposed) if (_connectionState == ConnectionState.Disposed)
{ {
_logger.HttpConnectionClosed(_connectionId); _logger.HttpConnectionClosed();
return; return;
} }
@ -163,7 +170,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
connectUrl = CreateConnectUrl(Url, negotiationResponse); connectUrl = CreateConnectUrl(Url, negotiationResponse);
} }
_logger.StartingTransport(_connectionId, _transport.GetType().Name, connectUrl); _logger.StartingTransport(_transport, connectUrl);
await StartTransport(connectUrl); await StartTransport(connectUrl);
} }
catch catch
@ -193,16 +200,17 @@ namespace Microsoft.AspNetCore.Sockets.Client
// to make sure that the message removed from the channel is processed before we drain the queue. // to make sure that the message removed from the channel is processed before we drain the queue.
// There is a short window between we start the channel and assign the _receiveLoopTask a value. // There is a short window between we start the channel and assign the _receiveLoopTask a value.
// To make sure that _receiveLoopTask can be awaited (i.e. is not null) we need to await _startTask. // To make sure that _receiveLoopTask can be awaited (i.e. is not null) we need to await _startTask.
_logger.ProcessRemainingMessages(_connectionId); _logger.ProcessRemainingMessages();
await _startTcs.Task; await _startTcs.Task;
await _receiveLoopTask; await _receiveLoopTask;
_logger.DrainEvents(_connectionId); _logger.DrainEvents();
await Task.WhenAny(_eventQueue.Drain().NoThrow(), Task.Delay(_eventQueueDrainTimeout)); await Task.WhenAny(_eventQueue.Drain().NoThrow(), Task.Delay(_eventQueueDrainTimeout));
_logger.CompleteClosed(_connectionId); _logger.CompleteClosed();
_logScope.ConnectionId = null;
// At this point the connection can be either in the Connected or Disposed state. The state should be changed // At this point the connection can be either in the Connected or Disposed state. The state should be changed
// to the Disconnected state only if it was in the Connected state. // to the Disconnected state only if it was in the Connected state.
@ -325,7 +333,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
// Start the transport, giving it one end of the pipeline // Start the transport, giving it one end of the pipeline
try try
{ {
await _transport.StartAsync(connectUrl, applicationSide, GetTransferMode(), _connectionId, this); await _transport.StartAsync(connectUrl, applicationSide, GetTransferMode(), this);
// actual transfer mode can differ from the one that was requested so set it on the feature // actual transfer mode can differ from the one that was requested so set it on the feature
if (!_transport.Mode.HasValue) if (!_transport.Mode.HasValue)
@ -337,7 +345,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.ErrorStartingTransport(_connectionId, _transport.GetType().Name, ex); _logger.ErrorStartingTransport(_transport, ex);
throw; throw;
} }
} }
@ -369,13 +377,13 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
try try
{ {
_logger.HttpReceiveStarted(_connectionId); _logger.HttpReceiveStarted();
while (await Input.WaitToReadAsync()) while (await Input.WaitToReadAsync())
{ {
if (_connectionState != ConnectionState.Connected) if (_connectionState != ConnectionState.Connected)
{ {
_logger.SkipRaisingReceiveEvent(_connectionId); _logger.SkipRaisingReceiveEvent();
// drain // drain
Input.TryRead(out _); Input.TryRead(out _);
continue; continue;
@ -383,10 +391,10 @@ namespace Microsoft.AspNetCore.Sockets.Client
if (Input.TryRead(out var buffer)) if (Input.TryRead(out var buffer))
{ {
_logger.ScheduleReceiveEvent(_connectionId); _logger.ScheduleReceiveEvent();
_ = _eventQueue.Enqueue(async () => _ = _eventQueue.Enqueue(async () =>
{ {
_logger.RaiseReceiveEvent(_connectionId); _logger.RaiseReceiveEvent();
// Copying the callbacks to avoid concurrency issues // Copying the callbacks to avoid concurrency issues
ReceiveCallback[] callbackCopies; ReceiveCallback[] callbackCopies;
@ -404,14 +412,14 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.ExceptionThrownFromCallback(_connectionId, nameof(OnReceived), ex); _logger.ExceptionThrownFromCallback(nameof(OnReceived), ex);
} }
} }
}); });
} }
else else
{ {
_logger.FailedReadingMessage(_connectionId); _logger.FailedReadingMessage();
} }
} }
@ -420,10 +428,10 @@ namespace Microsoft.AspNetCore.Sockets.Client
catch (Exception ex) catch (Exception ex)
{ {
Output.TryComplete(ex); Output.TryComplete(ex);
_logger.ErrorReceiving(_connectionId, ex); _logger.ErrorReceiving(ex);
} }
_logger.EndReceive(_connectionId); _logger.EndReceive();
} }
public async Task SendAsync(byte[] data, CancellationToken cancellationToken = default) => public async Task SendAsync(byte[] data, CancellationToken cancellationToken = default) =>
@ -449,7 +457,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
var sendTcs = new TaskCompletionSource<object>(TaskCreationOptions.RunContinuationsAsynchronously); var sendTcs = new TaskCompletionSource<object>(TaskCreationOptions.RunContinuationsAsynchronously);
var message = new SendMessage(data, sendTcs); var message = new SendMessage(data, sendTcs);
_logger.SendingMessage(_connectionId); _logger.SendingMessage();
while (await Output.WaitToWriteAsync(cancellationToken)) while (await Output.WaitToWriteAsync(cancellationToken))
{ {
@ -483,7 +491,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
if (!(_connectionState == ConnectionState.Connecting || _connectionState == ConnectionState.Connected)) if (!(_connectionState == ConnectionState.Connecting || _connectionState == ConnectionState.Connected))
{ {
_logger.SkippingStop(_connectionId); _logger.SkippingStop();
return; return;
} }
} }
@ -496,7 +504,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
// See comment at AbortAsync for more discussion on the thread-safety of this. // See comment at AbortAsync for more discussion on the thread-safety of this.
_abortException = exception; _abortException = exception;
_logger.StoppingClient(_connectionId); _logger.StoppingClient();
try try
{ {
@ -538,15 +546,16 @@ namespace Microsoft.AspNetCore.Sockets.Client
if (ChangeState(to: ConnectionState.Disposed) == ConnectionState.Disposed) if (ChangeState(to: ConnectionState.Disposed) == ConnectionState.Disposed)
{ {
_logger.SkippingDispose(_connectionId); _logger.SkippingDispose();
// the connection was already disposed // the connection was already disposed
return; return;
} }
_logger.DisposingClient(_connectionId); _logger.DisposingClient();
_httpClient?.Dispose(); _httpClient?.Dispose();
_scopeDisposable.Dispose();
} }
public IDisposable OnReceived(Func<byte[], object, Task> callback, object state) public IDisposable OnReceived(Func<byte[], object, Task> callback, object state)
@ -605,7 +614,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
_connectionState = to; _connectionState = to;
} }
_logger.ConnectionStateChanged(_connectionId, state, to); _logger.ConnectionStateChanged(state, to);
return state; return state;
} }
} }
@ -616,7 +625,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
var state = _connectionState; var state = _connectionState;
_connectionState = to; _connectionState = to;
_logger.ConnectionStateChanged(_connectionId, state, to); _logger.ConnectionStateChanged(state, to);
return state; return state;
} }
} }

View File

@ -9,7 +9,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
public interface ITransport public interface ITransport
{ {
Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, string connectionId, IConnection connection); Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, IConnection connection);
Task StopAsync(); Task StopAsync();
TransferMode? Mode { get; } TransferMode? Mode { get; }
} }

View File

@ -0,0 +1,73 @@
// 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.Collections;
using System.Collections.Generic;
using System.Globalization;
namespace Microsoft.AspNetCore.Sockets.Http.Internal
{
public class ConnectionLogScope : IReadOnlyList<KeyValuePair<string, object>>
{
private string _cachedToString;
private string _connectionId;
public string ConnectionId
{
get
{
return _connectionId;
}
set
{
_cachedToString = null;
_connectionId = value;
}
}
public KeyValuePair<string, object> this[int index]
{
get
{
if (Count == 1 && index == 0)
{
return new KeyValuePair<string, object>("ClientConnectionId", ConnectionId);
}
throw new ArgumentOutOfRangeException(nameof(index));
}
}
public int Count => string.IsNullOrEmpty(ConnectionId) ? 0 : 1;
public IEnumerator<KeyValuePair<string, object>> GetEnumerator()
{
for (var i = 0; i < Count; ++i)
{
yield return this[i];
}
}
IEnumerator IEnumerable.GetEnumerator()
{
return GetEnumerator();
}
public override string ToString()
{
if (_cachedToString == null)
{
if (!string.IsNullOrEmpty(ConnectionId))
{
_cachedToString = string.Format(
CultureInfo.InvariantCulture,
"ClientConnectionId:{0}",
ConnectionId);
}
}
return _cachedToString;
}
}
}

View File

@ -10,573 +10,429 @@ namespace Microsoft.AspNetCore.Sockets.Client.Internal
internal static class SocketClientLoggerExtensions internal static class SocketClientLoggerExtensions
{ {
// Category: Shared with LongPollingTransport, WebSocketsTransport and ServerSentEventsTransport // Category: Shared with LongPollingTransport, WebSocketsTransport and ServerSentEventsTransport
private static readonly Action<ILogger, DateTime, string, TransferMode, Exception> _startTransport = private static readonly Action<ILogger, TransferMode, Exception> _startTransport =
LoggerMessage.Define<DateTime, string, TransferMode>(LogLevel.Information, new EventId(1, nameof(StartTransport)), "{time}: Connection Id {connectionId}: Starting transport. Transfer mode: {transferMode}."); LoggerMessage.Define<TransferMode>(LogLevel.Information, new EventId(1, nameof(StartTransport)), "Starting transport. Transfer mode: {transferMode}.");
private static readonly Action<ILogger, DateTime, string, Exception> _transportStopped = private static readonly Action<ILogger, Exception> _transportStopped =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(2, nameof(TransportStopped)), "{time}: Connection Id {connectionId}: Transport stopped."); LoggerMessage.Define(LogLevel.Debug, new EventId(2, nameof(TransportStopped)), "Transport stopped.");
private static readonly Action<ILogger, DateTime, string, Exception> _startReceive = private static readonly Action<ILogger, Exception> _startReceive =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(3, nameof(StartReceive)), "{time}: Connection Id {connectionId}: Starting receive loop."); LoggerMessage.Define(LogLevel.Debug, new EventId(3, nameof(StartReceive)), "Starting receive loop.");
private static readonly Action<ILogger, DateTime, string, Exception> _receiveStopped = private static readonly Action<ILogger, Exception> _receiveStopped =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(4, nameof(ReceiveStopped)), "{time}: Connection Id {connectionId}: Receive loop stopped."); LoggerMessage.Define(LogLevel.Debug, new EventId(4, nameof(ReceiveStopped)), "Receive loop stopped.");
private static readonly Action<ILogger, DateTime, string, Exception> _receiveCanceled = private static readonly Action<ILogger, Exception> _receiveCanceled =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(5, nameof(ReceiveCanceled)), "{time}: Connection Id {connectionId}: Receive loop canceled."); LoggerMessage.Define(LogLevel.Debug, new EventId(5, nameof(ReceiveCanceled)), "Receive loop canceled.");
private static readonly Action<ILogger, DateTime, string, Exception> _transportStopping = private static readonly Action<ILogger, Exception> _transportStopping =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(6, nameof(TransportStopping)), "{time}: Connection Id {connectionId}: Transport is stopping."); LoggerMessage.Define(LogLevel.Information, new EventId(6, nameof(TransportStopping)), "Transport is stopping.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendStarted = private static readonly Action<ILogger, Exception> _sendStarted =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(7, nameof(SendStarted)), "{time}: Connection Id {connectionId}: Starting the send loop."); LoggerMessage.Define(LogLevel.Debug, new EventId(7, nameof(SendStarted)), "Starting the send loop.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendStopped = private static readonly Action<ILogger, Exception> _sendStopped =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(8, nameof(SendStopped)), "{time}: Connection Id {connectionId}: Send loop stopped."); LoggerMessage.Define(LogLevel.Debug, new EventId(8, nameof(SendStopped)), "Send loop stopped.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendCanceled = private static readonly Action<ILogger, Exception> _sendCanceled =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(9, nameof(SendCanceled)), "{time}: Connection Id {connectionId}: Send loop canceled."); LoggerMessage.Define(LogLevel.Debug, new EventId(9, nameof(SendCanceled)), "Send loop canceled.");
// Category: WebSocketsTransport // Category: WebSocketsTransport
private static readonly Action<ILogger, DateTime, string, WebSocketCloseStatus?, Exception> _webSocketClosed = private static readonly Action<ILogger, WebSocketCloseStatus?, Exception> _webSocketClosed =
LoggerMessage.Define<DateTime, string, WebSocketCloseStatus?>(LogLevel.Information, new EventId(10, nameof(WebSocketClosed)), "{time}: Connection Id {connectionId}: Websocket closed by the server. Close status {closeStatus}."); LoggerMessage.Define<WebSocketCloseStatus?>(LogLevel.Information, new EventId(10, nameof(WebSocketClosed)), "Websocket closed by the server. Close status {closeStatus}.");
private static readonly Action<ILogger, DateTime, string, WebSocketMessageType, int, bool, Exception> _messageReceived = private static readonly Action<ILogger, WebSocketMessageType, int, bool, Exception> _messageReceived =
LoggerMessage.Define<DateTime, string, WebSocketMessageType, int, bool>(LogLevel.Debug, new EventId(11, nameof(MessageReceived)), "{time}: Connection Id {connectionId}: Message received. Type: {messageType}, size: {count}, EndOfMessage: {endOfMessage}."); LoggerMessage.Define<WebSocketMessageType, int, bool>(LogLevel.Debug, new EventId(11, nameof(MessageReceived)), "Message received. Type: {messageType}, size: {count}, EndOfMessage: {endOfMessage}.");
private static readonly Action<ILogger, DateTime, string, int, Exception> _messageToApp = private static readonly Action<ILogger, int, Exception> _messageToApp =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(12, nameof(MessageToApp)), "{time}: Connection Id {connectionId}: Passing message to application. Payload size: {count}."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(12, nameof(MessageToApp)), "Passing message to application. Payload size: {count}.");
private static readonly Action<ILogger, DateTime, string, int, Exception> _receivedFromApp = private static readonly Action<ILogger, int, Exception> _receivedFromApp =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(13, nameof(ReceivedFromApp)), "{time}: Connection Id {connectionId}: Received message from application. Payload size: {count}."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(13, nameof(ReceivedFromApp)), "Received message from application. Payload size: {count}.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendMessageCanceled = private static readonly Action<ILogger, Exception> _sendMessageCanceled =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(14, nameof(SendMessageCanceled)), "{time}: Connection Id {connectionId}: Sending a message canceled."); LoggerMessage.Define(LogLevel.Information, new EventId(14, nameof(SendMessageCanceled)), "Sending a message canceled.");
private static readonly Action<ILogger, DateTime, string, Exception> _errorSendingMessage = private static readonly Action<ILogger, Exception> _errorSendingMessage =
LoggerMessage.Define<DateTime, string>(LogLevel.Error, new EventId(15, nameof(ErrorSendingMessage)), "{time}: Connection Id {connectionId}: Error while sending a message."); LoggerMessage.Define(LogLevel.Error, new EventId(15, nameof(ErrorSendingMessage)), "Error while sending a message.");
private static readonly Action<ILogger, DateTime, string, Exception> _closingWebSocket = private static readonly Action<ILogger, Exception> _closingWebSocket =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(16, nameof(ClosingWebSocket)), "{time}: Connection Id {connectionId}: Closing WebSocket."); LoggerMessage.Define(LogLevel.Information, new EventId(16, nameof(ClosingWebSocket)), "Closing WebSocket.");
private static readonly Action<ILogger, DateTime, string, Exception> _closingWebSocketFailed = private static readonly Action<ILogger, Exception> _closingWebSocketFailed =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(17, nameof(ClosingWebSocketFailed)), "{time}: Connection Id {connectionId}: Closing webSocket failed."); LoggerMessage.Define(LogLevel.Information, new EventId(17, nameof(ClosingWebSocketFailed)), "Closing webSocket failed.");
private static readonly Action<ILogger, DateTime, string, Exception> _cancelMessage = private static readonly Action<ILogger, Exception> _cancelMessage =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(18, nameof(CancelMessage)), "{time}: Connection Id {connectionId}: Canceled passing message to application."); LoggerMessage.Define(LogLevel.Debug, new EventId(18, nameof(CancelMessage)), "Canceled passing message to application.");
// Category: ServerSentEventsTransport and LongPollingTransport // Category: ServerSentEventsTransport and LongPollingTransport
private static readonly Action<ILogger, DateTime, string, int, Uri, Exception> _sendingMessages = private static readonly Action<ILogger, int, Uri, Exception> _sendingMessages =
LoggerMessage.Define<DateTime, string, int, Uri>(LogLevel.Debug, new EventId(10, nameof(SendingMessages)), "{time}: Connection Id {connectionId}: Sending {count} message(s) to the server using url: {url}."); LoggerMessage.Define<int, Uri>(LogLevel.Debug, new EventId(10, nameof(SendingMessages)), "Sending {count} message(s) to the server using url: {url}.");
private static readonly Action<ILogger, DateTime, string, Exception> _sentSuccessfully = private static readonly Action<ILogger, Exception> _sentSuccessfully =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(11, nameof(SentSuccessfully)), "{time}: Connection Id {connectionId}: Message(s) sent successfully."); LoggerMessage.Define(LogLevel.Debug, new EventId(11, nameof(SentSuccessfully)), "Message(s) sent successfully.");
private static readonly Action<ILogger, DateTime, string, Exception> _noMessages = private static readonly Action<ILogger, Exception> _noMessages =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(12, nameof(NoMessages)), "{time}: Connection Id {connectionId}: No messages in batch to send."); LoggerMessage.Define(LogLevel.Debug, new EventId(12, nameof(NoMessages)), "No messages in batch to send.");
private static readonly Action<ILogger, DateTime, string, Uri, Exception> _errorSending = private static readonly Action<ILogger, Uri, Exception> _errorSending =
LoggerMessage.Define<DateTime, string, Uri>(LogLevel.Error, new EventId(13, nameof(ErrorSending)), "{time}: Connection Id {connectionId}: Error while sending to '{url}'."); LoggerMessage.Define<Uri>(LogLevel.Error, new EventId(13, nameof(ErrorSending)), "Error while sending to '{url}'.");
// Category: ServerSentEventsTransport // Category: ServerSentEventsTransport
private static readonly Action<ILogger, DateTime, string, Exception> _eventStreamEnded = private static readonly Action<ILogger, Exception> _eventStreamEnded =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(14, nameof(EventStreamEnded)), "{time}: Connection Id {connectionId}: Server-Sent Event Stream ended."); LoggerMessage.Define(LogLevel.Debug, new EventId(14, nameof(EventStreamEnded)), "Server-Sent Event Stream ended.");
// Category: LongPollingTransport // Category: LongPollingTransport
private static readonly Action<ILogger, DateTime, string, Exception> _closingConnection = private static readonly Action<ILogger, Exception> _closingConnection =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(14, nameof(ClosingConnection)), "{time}: Connection Id {connectionId}: The server is closing the connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(14, nameof(ClosingConnection)), "The server is closing the connection.");
private static readonly Action<ILogger, DateTime, string, Exception> _receivedMessages = private static readonly Action<ILogger, Exception> _receivedMessages =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(15, nameof(ReceivedMessages)), "{time}: Connection Id {connectionId}: Received messages from the server."); LoggerMessage.Define(LogLevel.Debug, new EventId(15, nameof(ReceivedMessages)), "Received messages from the server.");
private static readonly Action<ILogger, DateTime, string, Uri, Exception> _errorPolling = private static readonly Action<ILogger, Uri, Exception> _errorPolling =
LoggerMessage.Define<DateTime, string, Uri>(LogLevel.Error, new EventId(16, nameof(ErrorPolling)), "{time}: Connection Id {connectionId}: Error while polling '{pollUrl}'."); LoggerMessage.Define<Uri>(LogLevel.Error, new EventId(16, nameof(ErrorPolling)), "Error while polling '{pollUrl}'.");
// Category: HttpConnection // Category: HttpConnection
private static readonly Action<ILogger, DateTime, Exception> _httpConnectionStarting = private static readonly Action<ILogger, Exception> _httpConnectionStarting =
LoggerMessage.Define<DateTime>(LogLevel.Debug, new EventId(1, nameof(HttpConnectionStarting)), "{time}: Starting connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(1, nameof(HttpConnectionStarting)), "Starting connection.");
private static readonly Action<ILogger, DateTime, string, Exception> _httpConnectionClosed = private static readonly Action<ILogger, Exception> _httpConnectionClosed =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(2, nameof(HttpConnectionClosed)), "{time}: Connection Id {connectionId}: Connection was closed from a different thread."); LoggerMessage.Define(LogLevel.Debug, new EventId(2, nameof(HttpConnectionClosed)), "Connection was closed from a different thread.");
private static readonly Action<ILogger, DateTime, string, string, Uri, Exception> _startingTransport = private static readonly Action<ILogger, string, Uri, Exception> _startingTransport =
LoggerMessage.Define<DateTime, string, string, Uri>(LogLevel.Debug, new EventId(3, nameof(StartingTransport)), "{time}: Connection Id {connectionId}: Starting transport '{transport}' with Url: {url}."); LoggerMessage.Define<string, Uri>(LogLevel.Debug, new EventId(3, nameof(StartingTransport)), "Starting transport '{transport}' with Url: {url}.");
private static readonly Action<ILogger, DateTime, string, Exception> _raiseConnected = private static readonly Action<ILogger, Exception> _raiseConnected =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(4, nameof(RaiseConnected)), "{time}: Connection Id {connectionId}: Raising Connected event."); LoggerMessage.Define(LogLevel.Debug, new EventId(4, nameof(RaiseConnected)), "Raising Connected event.");
private static readonly Action<ILogger, DateTime, string, Exception> _processRemainingMessages = private static readonly Action<ILogger, Exception> _processRemainingMessages =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(5, nameof(ProcessRemainingMessages)), "{time}: Connection Id {connectionId}: Ensuring all outstanding messages are processed."); LoggerMessage.Define(LogLevel.Debug, new EventId(5, nameof(ProcessRemainingMessages)), "Ensuring all outstanding messages are processed.");
private static readonly Action<ILogger, DateTime, string, Exception> _drainEvents = private static readonly Action<ILogger, Exception> _drainEvents =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(6, nameof(DrainEvents)), "{time}: Connection Id {connectionId}: Draining event queue."); LoggerMessage.Define(LogLevel.Debug, new EventId(6, nameof(DrainEvents)), "Draining event queue.");
private static readonly Action<ILogger, DateTime, string, Exception> _completeClosed = private static readonly Action<ILogger, Exception> _completeClosed =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(7, nameof(CompleteClosed)), "{time}: Connection Id {connectionId}: Completing Closed task."); LoggerMessage.Define(LogLevel.Debug, new EventId(7, nameof(CompleteClosed)), "Completing Closed task.");
private static readonly Action<ILogger, DateTime, Uri, Exception> _establishingConnection = private static readonly Action<ILogger, Uri, Exception> _establishingConnection =
LoggerMessage.Define<DateTime, Uri>(LogLevel.Debug, new EventId(8, nameof(EstablishingConnection)), "{time}: Establishing Connection at: {url}."); LoggerMessage.Define<Uri>(LogLevel.Debug, new EventId(8, nameof(EstablishingConnection)), "Establishing Connection at: {url}.");
private static readonly Action<ILogger, DateTime, Uri, Exception> _errorWithNegotiation = private static readonly Action<ILogger, Uri, Exception> _errorWithNegotiation =
LoggerMessage.Define<DateTime, Uri>(LogLevel.Error, new EventId(9, nameof(ErrorWithNegotiation)), "{time}: Failed to start connection. Error getting negotiation response from '{url}'."); LoggerMessage.Define<Uri>(LogLevel.Error, new EventId(9, nameof(ErrorWithNegotiation)), "Failed to start connection. Error getting negotiation response from '{url}'.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _errorStartingTransport = private static readonly Action<ILogger, string, Exception> _errorStartingTransport =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Error, new EventId(10, nameof(ErrorStartingTransport)), "{time}: Connection Id {connectionId}: Failed to start connection. Error starting transport '{transport}'."); LoggerMessage.Define<string>(LogLevel.Error, new EventId(10, nameof(ErrorStartingTransport)), "Failed to start connection. Error starting transport '{transport}'.");
private static readonly Action<ILogger, DateTime, string, Exception> _httpReceiveStarted = private static readonly Action<ILogger, Exception> _httpReceiveStarted =
LoggerMessage.Define<DateTime, string>(LogLevel.Trace, new EventId(11, nameof(HttpReceiveStarted)), "{time}: Connection Id {connectionId}: Beginning receive loop."); LoggerMessage.Define(LogLevel.Trace, new EventId(11, nameof(HttpReceiveStarted)), "Beginning receive loop.");
private static readonly Action<ILogger, DateTime, string, Exception> _skipRaisingReceiveEvent = private static readonly Action<ILogger, Exception> _skipRaisingReceiveEvent =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(12, nameof(SkipRaisingReceiveEvent)), "{time}: Connection Id {connectionId}: Message received but connection is not connected. Skipping raising Received event."); LoggerMessage.Define(LogLevel.Debug, new EventId(12, nameof(SkipRaisingReceiveEvent)), "Message received but connection is not connected. Skipping raising Received event.");
private static readonly Action<ILogger, DateTime, string, Exception> _scheduleReceiveEvent = private static readonly Action<ILogger, Exception> _scheduleReceiveEvent =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(13, nameof(ScheduleReceiveEvent)), "{time}: Connection Id {connectionId}: Scheduling raising Received event."); LoggerMessage.Define(LogLevel.Debug, new EventId(13, nameof(ScheduleReceiveEvent)), "Scheduling raising Received event.");
private static readonly Action<ILogger, DateTime, string, Exception> _raiseReceiveEvent = private static readonly Action<ILogger, Exception> _raiseReceiveEvent =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(14, nameof(RaiseReceiveEvent)), "{time}: Connection Id {connectionId}: Raising Received event."); LoggerMessage.Define(LogLevel.Debug, new EventId(14, nameof(RaiseReceiveEvent)), "Raising Received event.");
private static readonly Action<ILogger, DateTime, string, Exception> _failedReadingMessage = private static readonly Action<ILogger, Exception> _failedReadingMessage =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(15, nameof(FailedReadingMessage)), "{time}: Connection Id {connectionId}: Could not read message."); LoggerMessage.Define(LogLevel.Debug, new EventId(15, nameof(FailedReadingMessage)), "Could not read message.");
private static readonly Action<ILogger, DateTime, string, Exception> _errorReceiving = private static readonly Action<ILogger, Exception> _errorReceiving =
LoggerMessage.Define<DateTime, string>(LogLevel.Error, new EventId(16, nameof(ErrorReceiving)), "{time}: Connection Id {connectionId}: Error receiving message."); LoggerMessage.Define(LogLevel.Error, new EventId(16, nameof(ErrorReceiving)), "Error receiving message.");
private static readonly Action<ILogger, DateTime, string, Exception> _endReceive = private static readonly Action<ILogger, Exception> _endReceive =
LoggerMessage.Define<DateTime, string>(LogLevel.Trace, new EventId(17, nameof(EndReceive)), "{time}: Connection Id {connectionId}: Ending receive loop."); LoggerMessage.Define(LogLevel.Trace, new EventId(17, nameof(EndReceive)), "Ending receive loop.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendingMessage = private static readonly Action<ILogger, Exception> _sendingMessage =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(18, nameof(SendingMessage)), "{time}: Connection Id {connectionId}: Sending message."); LoggerMessage.Define(LogLevel.Debug, new EventId(18, nameof(SendingMessage)), "Sending message.");
private static readonly Action<ILogger, DateTime, string, Exception> _stoppingClient = private static readonly Action<ILogger, Exception> _stoppingClient =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(19, nameof(StoppingClient)), "{time}: Connection Id {connectionId}: Stopping client."); LoggerMessage.Define(LogLevel.Information, new EventId(19, nameof(StoppingClient)), "Stopping client.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _exceptionThrownFromCallback = private static readonly Action<ILogger, string, Exception> _exceptionThrownFromCallback =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Error, new EventId(20, nameof(ExceptionThrownFromCallback)), "{time}: Connection Id {connectionId}: An exception was thrown from the '{callback}' callback."); LoggerMessage.Define<string>(LogLevel.Error, new EventId(20, nameof(ExceptionThrownFromCallback)), "An exception was thrown from the '{callback}' callback.");
private static readonly Action<ILogger, DateTime, string, Exception> _disposingClient = private static readonly Action<ILogger, Exception> _disposingClient =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(21, nameof(DisposingClient)), "{time}: Connection Id {connectionId}: Disposing client."); LoggerMessage.Define(LogLevel.Information, new EventId(21, nameof(DisposingClient)), "Disposing client.");
private static readonly Action<ILogger, DateTime, string, Exception> _abortingClient = private static readonly Action<ILogger, Exception> _abortingClient =
LoggerMessage.Define<DateTime, string>(LogLevel.Error, new EventId(22, nameof(AbortingClient)), "{time}: Connection Id {connectionId}: Aborting client."); LoggerMessage.Define(LogLevel.Error, new EventId(22, nameof(AbortingClient)), "Aborting client.");
private static readonly Action<ILogger, Exception> _errorDuringClosedEvent = private static readonly Action<ILogger, Exception> _errorDuringClosedEvent =
LoggerMessage.Define(LogLevel.Error, new EventId(23, nameof(ErrorDuringClosedEvent)), "An exception was thrown in the handler for the Closed event."); LoggerMessage.Define(LogLevel.Error, new EventId(23, nameof(ErrorDuringClosedEvent)), "An exception was thrown in the handler for the Closed event.");
private static readonly Action<ILogger, DateTime, string, Exception> _skippingStop = private static readonly Action<ILogger, Exception> _skippingStop =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(24, nameof(SkippingStop)), "{time}: Connection Id {connectionId}: Skipping stop, connection is already stopped."); LoggerMessage.Define(LogLevel.Debug, new EventId(24, nameof(SkippingStop)), "Skipping stop, connection is already stopped.");
private static readonly Action<ILogger, DateTime, string, Exception> _skippingDispose = private static readonly Action<ILogger, Exception> _skippingDispose =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(25, nameof(SkippingDispose)), "{time}: Connection Id {connectionId}: Skipping dispose, connection is already disposed."); LoggerMessage.Define(LogLevel.Debug, new EventId(25, nameof(SkippingDispose)), "Skipping dispose, connection is already disposed.");
private static readonly Action<ILogger, DateTime, string, string, string, Exception> _connectionStateChanged = private static readonly Action<ILogger, string, string, Exception> _connectionStateChanged =
LoggerMessage.Define<DateTime, string, string, string>(LogLevel.Debug, new EventId(26, nameof(ConnectionStateChanged)), "{time}: Connection Id {connectionId}: Connection state changed from {previousState} to {newState}."); LoggerMessage.Define<string, string>(LogLevel.Debug, new EventId(26, nameof(ConnectionStateChanged)), "Connection state changed from {previousState} to {newState}.");
public static void StartTransport(this ILogger logger, string connectionId, TransferMode transferMode) public static void StartTransport(this ILogger logger, TransferMode transferMode)
{ {
if (logger.IsEnabled(LogLevel.Information)) _startTransport(logger, transferMode, null);
{
_startTransport(logger, DateTime.Now, connectionId, transferMode, null);
}
} }
public static void TransportStopped(this ILogger logger, string connectionId, Exception exception) public static void TransportStopped(this ILogger logger, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _transportStopped(logger, exception);
{
_transportStopped(logger, DateTime.Now, connectionId, exception);
}
} }
public static void StartReceive(this ILogger logger, string connectionId) public static void StartReceive(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _startReceive(logger, null);
{
_startReceive(logger, DateTime.Now, connectionId, null);
}
} }
public static void TransportStopping(this ILogger logger, string connectionId) public static void TransportStopping(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _transportStopping(logger, null);
{
_transportStopping(logger, DateTime.Now, connectionId, null);
}
} }
public static void WebSocketClosed(this ILogger logger, string connectionId, WebSocketCloseStatus? closeStatus) public static void WebSocketClosed(this ILogger logger, WebSocketCloseStatus? closeStatus)
{ {
if (logger.IsEnabled(LogLevel.Information)) _webSocketClosed(logger, closeStatus, null);
{
_webSocketClosed(logger, DateTime.Now, connectionId, closeStatus, null);
}
} }
public static void MessageReceived(this ILogger logger, string connectionId, WebSocketMessageType messageType, int count, bool endOfMessage) public static void MessageReceived(this ILogger logger, WebSocketMessageType messageType, int count, bool endOfMessage)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _messageReceived(logger, messageType, count, endOfMessage, null);
{
_messageReceived(logger, DateTime.Now, connectionId, messageType, count, endOfMessage, null);
}
} }
public static void MessageToApp(this ILogger logger, string connectionId, int count) public static void MessageToApp(this ILogger logger, int count)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _messageToApp(logger, count, null);
{
_messageToApp(logger, DateTime.Now, connectionId, count, null);
}
} }
public static void ReceiveCanceled(this ILogger logger, string connectionId) public static void ReceiveCanceled(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _receiveCanceled(logger, null);
{
_receiveCanceled(logger, DateTime.Now, connectionId, null);
}
} }
public static void ReceiveStopped(this ILogger logger, string connectionId) public static void ReceiveStopped(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _receiveStopped(logger, null);
{
_receiveStopped(logger, DateTime.Now, connectionId, null);
}
} }
public static void SendStarted(this ILogger logger, string connectionId) public static void SendStarted(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sendStarted(logger, null);
{
_sendStarted(logger, DateTime.Now, connectionId, null);
}
} }
public static void ReceivedFromApp(this ILogger logger, string connectionId, int count) public static void ReceivedFromApp(this ILogger logger, int count)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _receivedFromApp(logger, count, null);
{
_receivedFromApp(logger, DateTime.Now, connectionId, count, null);
}
} }
public static void SendMessageCanceled(this ILogger logger, string connectionId) public static void SendMessageCanceled(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _sendMessageCanceled(logger, null);
{
_sendMessageCanceled(logger, DateTime.Now, connectionId, null);
}
} }
public static void ErrorSendingMessage(this ILogger logger, string connectionId, Exception exception) public static void ErrorSendingMessage(this ILogger logger, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) _errorSendingMessage(logger, exception);
{
_errorSendingMessage(logger, DateTime.Now, connectionId, exception);
}
} }
public static void SendCanceled(this ILogger logger, string connectionId) public static void SendCanceled(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sendCanceled(logger, null);
{
_sendCanceled(logger, DateTime.Now, connectionId, null);
}
} }
public static void SendStopped(this ILogger logger, string connectionId) public static void SendStopped(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sendStopped(logger, null);
{
_sendStopped(logger, DateTime.Now, connectionId, null);
}
} }
public static void ClosingWebSocket(this ILogger logger, string connectionId) public static void ClosingWebSocket(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _closingWebSocket(logger, null);
{
_closingWebSocket(logger, DateTime.Now, connectionId, null);
}
} }
public static void ClosingWebSocketFailed(this ILogger logger, string connectionId, Exception exception) public static void ClosingWebSocketFailed(this ILogger logger, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Information)) _closingWebSocketFailed(logger, exception);
{
_closingWebSocketFailed(logger, DateTime.Now, connectionId, exception);
}
} }
public static void CancelMessage(this ILogger logger, string connectionId) public static void CancelMessage(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _cancelMessage(logger, null);
{
_cancelMessage(logger, DateTime.Now, connectionId, null);
}
} }
public static void SendingMessages(this ILogger logger, string connectionId, int count, Uri url) public static void SendingMessages(this ILogger logger, int count, Uri url)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sendingMessages(logger, count, url, null);
{
_sendingMessages(logger, DateTime.Now, connectionId, count, url, null);
}
} }
public static void SentSuccessfully(this ILogger logger, string connectionId) public static void SentSuccessfully(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sentSuccessfully(logger, null);
{
_sentSuccessfully(logger, DateTime.Now, connectionId, null);
}
} }
public static void NoMessages(this ILogger logger, string connectionId) public static void NoMessages(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _noMessages(logger, null);
{
_noMessages(logger, DateTime.Now, connectionId, null);
}
} }
public static void ErrorSending(this ILogger logger, string connectionId, Uri url, Exception exception) public static void ErrorSending(this ILogger logger, Uri url, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) _errorSending(logger, url, exception);
{
_errorSending(logger, DateTime.Now, connectionId, url, exception);
}
} }
public static void EventStreamEnded(this ILogger logger, string connectionId) public static void EventStreamEnded(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _eventStreamEnded(logger, null);
{
_eventStreamEnded(logger, DateTime.Now, connectionId, null);
}
} }
public static void ClosingConnection(this ILogger logger, string connectionId) public static void ClosingConnection(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _closingConnection(logger, null);
{
_closingConnection(logger, DateTime.Now, connectionId, null);
}
} }
public static void ReceivedMessages(this ILogger logger, string connectionId) public static void ReceivedMessages(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _receivedMessages(logger, null);
{
_receivedMessages(logger, DateTime.Now, connectionId, null);
}
} }
public static void ErrorPolling(this ILogger logger, string connectionId, Uri pollUrl, Exception exception) public static void ErrorPolling(this ILogger logger, Uri pollUrl, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) _errorPolling(logger, pollUrl, exception);
{
_errorPolling(logger, DateTime.Now, connectionId, pollUrl, exception);
}
} }
public static void HttpConnectionStarting(this ILogger logger) public static void HttpConnectionStarting(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _httpConnectionStarting(logger, null);
{
_httpConnectionStarting(logger, DateTime.Now, null);
}
} }
public static void HttpConnectionClosed(this ILogger logger, string connectionId) public static void HttpConnectionClosed(this ILogger logger)
{
_httpConnectionClosed(logger, null);
}
public static void StartingTransport(this ILogger logger, ITransport transport, Uri url)
{ {
if (logger.IsEnabled(LogLevel.Debug)) if (logger.IsEnabled(LogLevel.Debug))
{ {
_httpConnectionClosed(logger, DateTime.Now, connectionId, null); _startingTransport(logger, transport.GetType().Name, url, null);
} }
} }
public static void StartingTransport(this ILogger logger, string connectionId, string transport, Uri url) public static void RaiseConnected(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _raiseConnected(logger, null);
{
_startingTransport(logger, DateTime.Now, connectionId, transport, url, null);
}
} }
public static void RaiseConnected(this ILogger logger, string connectionId) public static void ProcessRemainingMessages(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _processRemainingMessages(logger, null);
{
_raiseConnected(logger, DateTime.Now, connectionId, null);
}
} }
public static void ProcessRemainingMessages(this ILogger logger, string connectionId) public static void DrainEvents(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _drainEvents(logger, null);
{
_processRemainingMessages(logger, DateTime.Now, connectionId, null);
}
} }
public static void DrainEvents(this ILogger logger, string connectionId) public static void CompleteClosed(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _completeClosed(logger, null);
{
_drainEvents(logger, DateTime.Now, connectionId, null);
}
}
public static void CompleteClosed(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_completeClosed(logger, DateTime.Now, connectionId, null);
}
} }
public static void EstablishingConnection(this ILogger logger, Uri url) public static void EstablishingConnection(this ILogger logger, Uri url)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _establishingConnection(logger, url, null);
{
_establishingConnection(logger, DateTime.Now, url, null);
}
} }
public static void ErrorWithNegotiation(this ILogger logger, Uri url, Exception exception) public static void ErrorWithNegotiation(this ILogger logger, Uri url, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) _errorWithNegotiation(logger, url, exception);
{
_errorWithNegotiation(logger, DateTime.Now, url, exception);
}
} }
public static void ErrorStartingTransport(this ILogger logger, string connectionId, string transport, Exception exception) public static void ErrorStartingTransport(this ILogger logger, ITransport transport, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) if (logger.IsEnabled(LogLevel.Error))
{ {
_errorStartingTransport(logger, DateTime.Now, connectionId, transport, exception); _errorStartingTransport(logger, transport.GetType().Name, exception);
} }
} }
public static void HttpReceiveStarted(this ILogger logger, string connectionId) public static void HttpReceiveStarted(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _httpReceiveStarted(logger, null);
{
_httpReceiveStarted(logger, DateTime.Now, connectionId, null);
}
} }
public static void SkipRaisingReceiveEvent(this ILogger logger, string connectionId) public static void SkipRaisingReceiveEvent(this ILogger logger)
{
_skipRaisingReceiveEvent(logger, null);
}
public static void ScheduleReceiveEvent(this ILogger logger)
{
_scheduleReceiveEvent(logger, null);
}
public static void RaiseReceiveEvent(this ILogger logger)
{
_raiseReceiveEvent(logger, null);
}
public static void FailedReadingMessage(this ILogger logger)
{
_failedReadingMessage(logger, null);
}
public static void ErrorReceiving(this ILogger logger, Exception exception)
{
_errorReceiving(logger, exception);
}
public static void EndReceive(this ILogger logger)
{
_endReceive(logger, null);
}
public static void SendingMessage(this ILogger logger)
{
_sendingMessage(logger, null);
}
public static void AbortingClient(this ILogger logger, Exception ex)
{
_abortingClient(logger, ex);
}
public static void StoppingClient(this ILogger logger)
{
_stoppingClient(logger, null);
}
public static void DisposingClient(this ILogger logger)
{
_disposingClient(logger, null);
}
public static void SkippingDispose(this ILogger logger)
{
_skippingDispose(logger, null);
}
public static void ConnectionStateChanged(this ILogger logger, HttpConnection.ConnectionState previousState, HttpConnection.ConnectionState newState)
{ {
if (logger.IsEnabled(LogLevel.Debug)) if (logger.IsEnabled(LogLevel.Debug))
{ {
_skipRaisingReceiveEvent(logger, DateTime.Now, connectionId, null); _connectionStateChanged(logger, previousState.ToString(), newState.ToString(), null);
} }
} }
public static void ScheduleReceiveEvent(this ILogger logger, string connectionId) public static void SkippingStop(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _skippingStop(logger, null);
{
_scheduleReceiveEvent(logger, DateTime.Now, connectionId, null);
}
} }
public static void RaiseReceiveEvent(this ILogger logger, string connectionId) public static void ExceptionThrownFromCallback(this ILogger logger, string callbackName, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _exceptionThrownFromCallback(logger, callbackName, exception);
{
_raiseReceiveEvent(logger, DateTime.Now, connectionId, null);
}
}
public static void FailedReadingMessage(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_failedReadingMessage(logger, DateTime.Now, connectionId, null);
}
}
public static void ErrorReceiving(this ILogger logger, string connectionId, Exception exception)
{
if (logger.IsEnabled(LogLevel.Error))
{
_errorReceiving(logger, DateTime.Now, connectionId, exception);
}
}
public static void EndReceive(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Trace))
{
_endReceive(logger, DateTime.Now, connectionId, null);
}
}
public static void SendingMessage(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_sendingMessage(logger, DateTime.Now, connectionId, null);
}
}
public static void AbortingClient(this ILogger logger, string connectionId, Exception ex)
{
if (logger.IsEnabled(LogLevel.Error))
{
_abortingClient(logger, DateTime.Now, connectionId, ex);
}
}
public static void StoppingClient(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Information))
{
_stoppingClient(logger, DateTime.Now, connectionId, null);
}
}
public static void DisposingClient(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Information))
{
_disposingClient(logger, DateTime.Now, connectionId, null);
}
}
public static void SkippingDispose(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_skippingDispose(logger, DateTime.Now, connectionId, null);
}
}
public static void ConnectionStateChanged(this ILogger logger, string connectionId, HttpConnection.ConnectionState previousState, HttpConnection.ConnectionState newState)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_connectionStateChanged(logger, DateTime.Now, connectionId, previousState.ToString(), newState.ToString(), null);
}
}
public static void SkippingStop(this ILogger logger, string connectionId)
{
if (logger.IsEnabled(LogLevel.Debug))
{
_skippingStop(logger, DateTime.Now, connectionId, null);
}
}
public static void ExceptionThrownFromCallback(this ILogger logger, string connectionId, string callbackName, Exception exception)
{
if (logger.IsEnabled(LogLevel.Error))
{
_exceptionThrownFromCallback(logger, DateTime.Now, connectionId, callbackName, exception);
}
} }
public static void ErrorDuringClosedEvent(this ILogger logger, Exception exception) public static void ErrorDuringClosedEvent(this ILogger logger, Exception exception)

View File

@ -23,7 +23,6 @@ namespace Microsoft.AspNetCore.Sockets.Client
private Channel<byte[], SendMessage> _application; private Channel<byte[], SendMessage> _application;
private Task _sender; private Task _sender;
private Task _poller; private Task _poller;
private string _connectionId;
private readonly CancellationTokenSource _transportCts = new CancellationTokenSource(); private readonly CancellationTokenSource _transportCts = new CancellationTokenSource();
@ -42,7 +41,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
_logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<LongPollingTransport>(); _logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<LongPollingTransport>();
} }
public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, string connectionId, IConnection connection) public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, IConnection connection)
{ {
if (requestedTransferMode != TransferMode.Binary && requestedTransferMode != TransferMode.Text) if (requestedTransferMode != TransferMode.Binary && requestedTransferMode != TransferMode.Text)
{ {
@ -53,17 +52,16 @@ namespace Microsoft.AspNetCore.Sockets.Client
_application = application; _application = application;
Mode = requestedTransferMode; Mode = requestedTransferMode;
_connectionId = connectionId;
_logger.StartTransport(_connectionId, Mode.Value); _logger.StartTransport(Mode.Value);
// Start sending and polling (ask for binary if the server supports it) // Start sending and polling (ask for binary if the server supports it)
_poller = Poll(url, _transportCts.Token); _poller = Poll(url, _transportCts.Token);
_sender = SendUtils.SendMessages(url, _application, _httpClient, _httpOptions, _transportCts, _logger, _connectionId); _sender = SendUtils.SendMessages(url, _application, _httpClient, _httpOptions, _transportCts, _logger);
Running = Task.WhenAll(_sender, _poller).ContinueWith(t => Running = Task.WhenAll(_sender, _poller).ContinueWith(t =>
{ {
_logger.TransportStopped(_connectionId, t.Exception?.InnerException); _logger.TransportStopped(t.Exception?.InnerException);
_application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null); _application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null);
return t; return t;
}).Unwrap(); }).Unwrap();
@ -73,7 +71,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
public async Task StopAsync() public async Task StopAsync()
{ {
_logger.TransportStopping(_connectionId); _logger.TransportStopping();
_transportCts.Cancel(); _transportCts.Cancel();
@ -89,7 +87,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
private async Task Poll(Uri pollUrl, CancellationToken cancellationToken) private async Task Poll(Uri pollUrl, CancellationToken cancellationToken)
{ {
_logger.StartReceive(_connectionId); _logger.StartReceive();
try try
{ {
while (!cancellationToken.IsCancellationRequested) while (!cancellationToken.IsCancellationRequested)
@ -115,14 +113,14 @@ namespace Microsoft.AspNetCore.Sockets.Client
if (response.StatusCode == HttpStatusCode.NoContent || cancellationToken.IsCancellationRequested) if (response.StatusCode == HttpStatusCode.NoContent || cancellationToken.IsCancellationRequested)
{ {
_logger.ClosingConnection(_connectionId); _logger.ClosingConnection();
// Transport closed or polling stopped, we're done // Transport closed or polling stopped, we're done
break; break;
} }
else else
{ {
_logger.ReceivedMessages(_connectionId); _logger.ReceivedMessages();
// Until Pipeline starts natively supporting BytesReader, this is the easiest way to do this. // Until Pipeline starts natively supporting BytesReader, this is the easiest way to do this.
var payload = await response.Content.ReadAsByteArrayAsync(); var payload = await response.Content.ReadAsByteArrayAsync();
@ -142,18 +140,18 @@ namespace Microsoft.AspNetCore.Sockets.Client
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
// transport is being closed // transport is being closed
_logger.ReceiveCanceled(_connectionId); _logger.ReceiveCanceled();
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.ErrorPolling(_connectionId, pollUrl, ex); _logger.ErrorPolling(pollUrl, ex);
throw; throw;
} }
finally finally
{ {
// Make sure the send loop is terminated // Make sure the send loop is terminated
_transportCts.Cancel(); _transportCts.Cancel();
_logger.ReceiveStopped(_connectionId); _logger.ReceiveStopped();
} }
} }
} }

View File

@ -17,9 +17,9 @@ namespace Microsoft.AspNetCore.Sockets.Client
internal static class SendUtils internal static class SendUtils
{ {
public static async Task SendMessages(Uri sendUrl, Channel<byte[], SendMessage> application, HttpClient httpClient, public static async Task SendMessages(Uri sendUrl, Channel<byte[], SendMessage> application, HttpClient httpClient,
HttpOptions httpOptions, CancellationTokenSource transportCts, ILogger logger, string connectionId) HttpOptions httpOptions, CancellationTokenSource transportCts, ILogger logger)
{ {
logger.SendStarted(connectionId); logger.SendStarted();
IList<SendMessage> messages = null; IList<SendMessage> messages = null;
try try
{ {
@ -35,7 +35,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
transportCts.Token.ThrowIfCancellationRequested(); transportCts.Token.ThrowIfCancellationRequested();
if (messages.Count > 0) if (messages.Count > 0)
{ {
logger.SendingMessages(connectionId, messages.Count, sendUrl); logger.SendingMessages(messages.Count, sendUrl);
// Send them in a single post // Send them in a single post
var request = new HttpRequestMessage(HttpMethod.Post, sendUrl); var request = new HttpRequestMessage(HttpMethod.Post, sendUrl);
@ -61,7 +61,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
var response = await httpClient.SendAsync(request, transportCts.Token); var response = await httpClient.SendAsync(request, transportCts.Token);
response.EnsureSuccessStatusCode(); response.EnsureSuccessStatusCode();
logger.SentSuccessfully(connectionId); logger.SentSuccessfully();
foreach (var message in messages) foreach (var message in messages)
{ {
message.SendResult?.TrySetResult(null); message.SendResult?.TrySetResult(null);
@ -69,7 +69,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
else else
{ {
logger.NoMessages(connectionId); logger.NoMessages();
} }
} }
} }
@ -84,11 +84,11 @@ namespace Microsoft.AspNetCore.Sockets.Client
message.SendResult?.TrySetCanceled(); message.SendResult?.TrySetCanceled();
} }
} }
logger.SendCanceled(connectionId); logger.SendCanceled();
} }
catch (Exception ex) catch (Exception ex)
{ {
logger.ErrorSending(connectionId, sendUrl, ex); logger.ErrorSending(sendUrl, ex);
if (messages != null) if (messages != null)
{ {
foreach (var message in messages) foreach (var message in messages)
@ -105,7 +105,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
transportCts.Cancel(); transportCts.Cancel();
} }
logger.SendStopped(connectionId); logger.SendStopped();
} }
public static void PrepareHttpRequest(HttpRequestMessage request, HttpOptions httpOptions) public static void PrepareHttpRequest(HttpRequestMessage request, HttpOptions httpOptions)

View File

@ -25,7 +25,6 @@ namespace Microsoft.AspNetCore.Sockets.Client
private readonly ILogger _logger; private readonly ILogger _logger;
private readonly CancellationTokenSource _transportCts = new CancellationTokenSource(); private readonly CancellationTokenSource _transportCts = new CancellationTokenSource();
private readonly ServerSentEventsMessageParser _parser = new ServerSentEventsMessageParser(); private readonly ServerSentEventsMessageParser _parser = new ServerSentEventsMessageParser();
private string _connectionId;
private Channel<byte[], SendMessage> _application; private Channel<byte[], SendMessage> _application;
@ -49,7 +48,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
_logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<ServerSentEventsTransport>(); _logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<ServerSentEventsTransport>();
} }
public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, string connectionId, IConnection connection) public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, IConnection connection)
{ {
if (requestedTransferMode != TransferMode.Binary && requestedTransferMode != TransferMode.Text) if (requestedTransferMode != TransferMode.Binary && requestedTransferMode != TransferMode.Text)
{ {
@ -58,16 +57,15 @@ namespace Microsoft.AspNetCore.Sockets.Client
_application = application; _application = application;
Mode = TransferMode.Text; // Server Sent Events is a text only transport Mode = TransferMode.Text; // Server Sent Events is a text only transport
_connectionId = connectionId;
_logger.StartTransport(_connectionId, Mode.Value); _logger.StartTransport(Mode.Value);
var sendTask = SendUtils.SendMessages(url, _application, _httpClient, _httpOptions, _transportCts, _logger, _connectionId); var sendTask = SendUtils.SendMessages(url, _application, _httpClient, _httpOptions, _transportCts, _logger);
var receiveTask = OpenConnection(_application, url, _transportCts.Token); var receiveTask = OpenConnection(_application, url, _transportCts.Token);
Running = Task.WhenAll(sendTask, receiveTask).ContinueWith(t => Running = Task.WhenAll(sendTask, receiveTask).ContinueWith(t =>
{ {
_logger.TransportStopped(_connectionId, t.Exception?.InnerException); _logger.TransportStopped(t.Exception?.InnerException);
_application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null); _application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null);
return t; return t;
@ -78,7 +76,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
private async Task OpenConnection(Channel<byte[], SendMessage> application, Uri url, CancellationToken cancellationToken) private async Task OpenConnection(Channel<byte[], SendMessage> application, Uri url, CancellationToken cancellationToken)
{ {
_logger.StartReceive(_connectionId); _logger.StartReceive();
var request = new HttpRequestMessage(HttpMethod.Get, url); var request = new HttpRequestMessage(HttpMethod.Get, url);
SendUtils.PrepareHttpRequest(request, _httpOptions); SendUtils.PrepareHttpRequest(request, _httpOptions);
@ -97,7 +95,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
var input = result.Buffer; var input = result.Buffer;
if (result.IsCancelled || (input.IsEmpty && result.IsCompleted)) if (result.IsCancelled || (input.IsEmpty && result.IsCompleted))
{ {
_logger.EventStreamEnded(_connectionId); _logger.EventStreamEnded();
break; break;
} }
@ -130,20 +128,20 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
_logger.ReceiveCanceled(_connectionId); _logger.ReceiveCanceled();
} }
finally finally
{ {
readCancellationRegistration.Dispose(); readCancellationRegistration.Dispose();
_transportCts.Cancel(); _transportCts.Cancel();
stream.Dispose(); stream.Dispose();
_logger.ReceiveStopped(_connectionId); _logger.ReceiveStopped();
} }
} }
public async Task StopAsync() public async Task StopAsync()
{ {
_logger.TransportStopping(_connectionId); _logger.TransportStopping();
_transportCts.Cancel(); _transportCts.Cancel();
_application.Writer.TryComplete(); _application.Writer.TryComplete();

View File

@ -22,7 +22,6 @@ namespace Microsoft.AspNetCore.Sockets.Client
private readonly CancellationTokenSource _transportCts = new CancellationTokenSource(); private readonly CancellationTokenSource _transportCts = new CancellationTokenSource();
private readonly CancellationTokenSource _receiveCts = new CancellationTokenSource(); private readonly CancellationTokenSource _receiveCts = new CancellationTokenSource();
private readonly ILogger _logger; private readonly ILogger _logger;
private string _connectionId;
public Task Running { get; private set; } = Task.CompletedTask; public Task Running { get; private set; } = Task.CompletedTask;
@ -55,7 +54,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
_logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<WebSocketsTransport>(); _logger = (loggerFactory ?? NullLoggerFactory.Instance).CreateLogger<WebSocketsTransport>();
} }
public async Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, string connectionId, IConnection connection) public async Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, IConnection connection)
{ {
if (url == null) if (url == null)
{ {
@ -74,9 +73,8 @@ namespace Microsoft.AspNetCore.Sockets.Client
_application = application; _application = application;
Mode = requestedTransferMode; Mode = requestedTransferMode;
_connectionId = connectionId;
_logger.StartTransport(_connectionId, Mode.Value); _logger.StartTransport(Mode.Value);
await Connect(url); await Connect(url);
var sendTask = SendMessages(url); var sendTask = SendMessages(url);
@ -87,7 +85,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
Running = Task.WhenAll(sendTask, receiveTask).ContinueWith(t => Running = Task.WhenAll(sendTask, receiveTask).ContinueWith(t =>
{ {
_webSocket.Dispose(); _webSocket.Dispose();
_logger.TransportStopped(_connectionId, t.Exception?.InnerException); _logger.TransportStopped(t.Exception?.InnerException);
_application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null); _application.Writer.TryComplete(t.IsFaulted ? t.Exception.InnerException : null);
return t; return t;
}).Unwrap(); }).Unwrap();
@ -95,7 +93,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
private async Task ReceiveMessages(Uri pollUrl) private async Task ReceiveMessages(Uri pollUrl)
{ {
_logger.StartReceive(_connectionId); _logger.StartReceive();
try try
{ {
@ -113,7 +111,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
receiveResult = await _webSocket.ReceiveAsync(buffer, _receiveCts.Token); receiveResult = await _webSocket.ReceiveAsync(buffer, _receiveCts.Token);
if (receiveResult.MessageType == WebSocketMessageType.Close) if (receiveResult.MessageType == WebSocketMessageType.Close)
{ {
_logger.WebSocketClosed(_connectionId, receiveResult.CloseStatus); _logger.WebSocketClosed(receiveResult.CloseStatus);
_application.Writer.Complete( _application.Writer.Complete(
receiveResult.CloseStatus == WebSocketCloseStatus.NormalClosure receiveResult.CloseStatus == WebSocketCloseStatus.NormalClosure
@ -123,7 +121,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
return; return;
} }
_logger.MessageReceived(_connectionId, receiveResult.MessageType, receiveResult.Count, receiveResult.EndOfMessage); _logger.MessageReceived(receiveResult.MessageType, receiveResult.Count, receiveResult.EndOfMessage);
var truncBuffer = new ArraySegment<byte>(buffer.Array, 0, receiveResult.Count); var truncBuffer = new ArraySegment<byte>(buffer.Array, 0, receiveResult.Count);
incomingMessage.Add(truncBuffer); incomingMessage.Add(truncBuffer);
@ -152,7 +150,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
if (!_transportCts.Token.IsCancellationRequested) if (!_transportCts.Token.IsCancellationRequested)
{ {
_logger.MessageToApp(_connectionId, messageBuffer.Length); _logger.MessageToApp(messageBuffer.Length);
while (await _application.Writer.WaitToWriteAsync(_transportCts.Token)) while (await _application.Writer.WaitToWriteAsync(_transportCts.Token))
{ {
if (_application.Writer.TryWrite(messageBuffer)) if (_application.Writer.TryWrite(messageBuffer))
@ -165,24 +163,24 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
_logger.CancelMessage(_connectionId); _logger.CancelMessage();
} }
} }
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
_logger.ReceiveCanceled(_connectionId); _logger.ReceiveCanceled();
} }
finally finally
{ {
_logger.ReceiveStopped(_connectionId); _logger.ReceiveStopped();
_transportCts.Cancel(); _transportCts.Cancel();
} }
} }
private async Task SendMessages(Uri sendUrl) private async Task SendMessages(Uri sendUrl)
{ {
_logger.SendStarted(_connectionId); _logger.SendStarted();
var webSocketMessageType = var webSocketMessageType =
Mode == TransferMode.Binary Mode == TransferMode.Binary
@ -197,7 +195,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
try try
{ {
_logger.ReceivedFromApp(_connectionId, message.Payload.Length); _logger.ReceivedFromApp(message.Payload.Length);
await _webSocket.SendAsync(new ArraySegment<byte>(message.Payload), webSocketMessageType, true, _transportCts.Token); await _webSocket.SendAsync(new ArraySegment<byte>(message.Payload), webSocketMessageType, true, _transportCts.Token);
@ -205,14 +203,14 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
_logger.SendMessageCanceled(_connectionId); _logger.SendMessageCanceled();
message.SendResult.SetCanceled(); message.SendResult.SetCanceled();
await CloseWebSocket(); await CloseWebSocket();
break; break;
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.ErrorSendingMessage(_connectionId, ex); _logger.ErrorSendingMessage(ex);
message.SendResult.SetException(ex); message.SendResult.SetException(ex);
await CloseWebSocket(); await CloseWebSocket();
throw; throw;
@ -222,11 +220,11 @@ namespace Microsoft.AspNetCore.Sockets.Client
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
_logger.SendCanceled(_connectionId); _logger.SendCanceled();
} }
finally finally
{ {
_logger.SendStopped(_connectionId); _logger.SendStopped();
TriggerCancel(); TriggerCancel();
} }
} }
@ -248,7 +246,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
public async Task StopAsync() public async Task StopAsync()
{ {
_logger.TransportStopping(_connectionId); _logger.TransportStopping();
await CloseWebSocket(); await CloseWebSocket();
@ -270,7 +268,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
// while the webSocket is being closed due to an error. // while the webSocket is being closed due to an error.
if (_webSocket.State != WebSocketState.Closed) if (_webSocket.State != WebSocketState.Closed)
{ {
_logger.ClosingWebSocket(_connectionId); _logger.ClosingWebSocket();
// We intentionally don't pass _transportCts.Token to CloseOutputAsync. The token can be cancelled // We intentionally don't pass _transportCts.Token to CloseOutputAsync. The token can be cancelled
// for reasons not related to webSocket in which case we would not close the websocket gracefully. // for reasons not related to webSocket in which case we would not close the websocket gracefully.
@ -284,7 +282,7 @@ namespace Microsoft.AspNetCore.Sockets.Client
{ {
// This is benign - the exception can happen due to the race described above because we would // This is benign - the exception can happen due to the race described above because we would
// try closing the webSocket twice. // try closing the webSocket twice.
_logger.ClosingWebSocketFailed(_connectionId, ex); _logger.ClosingWebSocketFailed(ex);
} }
} }

View File

@ -109,7 +109,7 @@ namespace Microsoft.AspNetCore.Sockets
return; return;
} }
_logger.EstablishedConnection(connection.ConnectionId, context.TraceIdentifier); _logger.EstablishedConnection();
// ServerSentEvents is a text protocol only // ServerSentEvents is a text protocol only
connection.TransportCapabilities = TransferMode.Text; connection.TransportCapabilities = TransferMode.Text;
@ -135,7 +135,7 @@ namespace Microsoft.AspNetCore.Sockets
return; return;
} }
_logger.EstablishedConnection(connection.ConnectionId, context.TraceIdentifier); _logger.EstablishedConnection();
var ws = new WebSocketsTransport(options.WebSockets, connection.Application, connection, _loggerFactory); var ws = new WebSocketsTransport(options.WebSockets, connection.Application, connection, _loggerFactory);
@ -196,7 +196,7 @@ namespace Microsoft.AspNetCore.Sockets
// Raise OnConnected for new connections only since polls happen all the time // Raise OnConnected for new connections only since polls happen all the time
if (connection.ApplicationTask == null) if (connection.ApplicationTask == null)
{ {
_logger.EstablishedConnection(connection.ConnectionId, connection.GetHttpContext().TraceIdentifier); _logger.EstablishedConnection();
connection.Metadata[ConnectionMetadataNames.Transport] = TransportType.LongPolling; connection.Metadata[ConnectionMetadataNames.Transport] = TransportType.LongPolling;
@ -204,7 +204,7 @@ namespace Microsoft.AspNetCore.Sockets
} }
else else
{ {
_logger.ResumingConnection(connection.ConnectionId, connection.GetHttpContext().TraceIdentifier); _logger.ResumingConnection();
} }
// REVIEW: Performance of this isn't great as this does a bunch of per request allocations // REVIEW: Performance of this isn't great as this does a bunch of per request allocations
@ -370,7 +370,7 @@ namespace Microsoft.AspNetCore.Sockets
// Get the bytes for the connection id // Get the bytes for the connection id
var negotiateResponseBuffer = Encoding.UTF8.GetBytes(GetNegotiatePayload(connection.ConnectionId, options)); var negotiateResponseBuffer = Encoding.UTF8.GetBytes(GetNegotiatePayload(connection.ConnectionId, options));
_logger.NegotiationRequest(connection.ConnectionId); _logger.NegotiationRequest();
// Write it out to the response with the right content length // Write it out to the response with the right content length
context.Response.ContentLength = negotiateResponseBuffer.Length; context.Response.ContentLength = negotiateResponseBuffer.Length;
@ -422,7 +422,7 @@ namespace Microsoft.AspNetCore.Sockets
var transport = (TransportType?)connection.Metadata[ConnectionMetadataNames.Transport]; var transport = (TransportType?)connection.Metadata[ConnectionMetadataNames.Transport];
if (transport == TransportType.WebSockets) if (transport == TransportType.WebSockets)
{ {
_logger.PostNotAllowedForWebSockets(connection.ConnectionId); _logger.PostNotAllowedForWebSockets();
context.Response.StatusCode = StatusCodes.Status405MethodNotAllowed; context.Response.StatusCode = StatusCodes.Status405MethodNotAllowed;
await context.Response.WriteAsync("POST requests are not allowed for WebSocket connections."); await context.Response.WriteAsync("POST requests are not allowed for WebSocket connections.");
return; return;
@ -438,7 +438,7 @@ namespace Microsoft.AspNetCore.Sockets
buffer = stream.ToArray(); buffer = stream.ToArray();
} }
_logger.ReceivedBytes(connection.ConnectionId, buffer.Length); _logger.ReceivedBytes(buffer.Length);
while (!connection.Application.Writer.TryWrite(buffer)) while (!connection.Application.Writer.TryWrite(buffer))
{ {
if (!await connection.Application.Writer.WaitToWriteAsync()) if (!await connection.Application.Writer.WaitToWriteAsync())
@ -454,7 +454,7 @@ namespace Microsoft.AspNetCore.Sockets
{ {
context.Response.ContentType = "text/plain"; context.Response.ContentType = "text/plain";
context.Response.StatusCode = StatusCodes.Status404NotFound; context.Response.StatusCode = StatusCodes.Status404NotFound;
_logger.TransportNotSupported(connection.ConnectionId, transportType); _logger.TransportNotSupported(transportType);
await context.Response.WriteAsync($"{transportType} transport not supported by this end point type"); await context.Response.WriteAsync($"{transportType} transport not supported by this end point type");
return false; return false;
} }
@ -469,7 +469,7 @@ namespace Microsoft.AspNetCore.Sockets
{ {
context.Response.ContentType = "text/plain"; context.Response.ContentType = "text/plain";
context.Response.StatusCode = StatusCodes.Status400BadRequest; context.Response.StatusCode = StatusCodes.Status400BadRequest;
_logger.CannotChangeTransport(connection.ConnectionId, transport.Value, transportType); _logger.CannotChangeTransport(transport.Value, transportType);
await context.Response.WriteAsync("Cannot change transports mid-connection"); await context.Response.WriteAsync("Cannot change transports mid-connection");
return false; return false;
} }

View File

@ -10,326 +10,239 @@ namespace Microsoft.AspNetCore.Sockets.Internal
internal static class SocketHttpLoggerExtensions internal static class SocketHttpLoggerExtensions
{ {
// Category: LongPollingTransport // Category: LongPollingTransport
private static readonly Action<ILogger, DateTime, string, string, Exception> _longPolling204 = private static readonly Action<ILogger, Exception> _longPolling204 =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Information, new EventId(1, nameof(LongPolling204)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Terminating Long Polling connection by sending 204 response."); LoggerMessage.Define(LogLevel.Information, new EventId(1, nameof(LongPolling204)), "Terminating Long Polling connection by sending 204 response.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _pollTimedOut = private static readonly Action<ILogger, Exception> _pollTimedOut =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Information, new EventId(2, nameof(PollTimedOut)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Poll request timed out. Sending 200 response to connection."); LoggerMessage.Define(LogLevel.Information, new EventId(2, nameof(PollTimedOut)), "Poll request timed out. Sending 200 response to connection.");
private static readonly Action<ILogger, DateTime, string, string, int, Exception> _longPollingWritingMessage = private static readonly Action<ILogger, int, Exception> _longPollingWritingMessage =
LoggerMessage.Define<DateTime, string, string, int>(LogLevel.Debug, new EventId(3, nameof(LongPollingWritingMessage)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Writing a {count} byte message to connection."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(3, nameof(LongPollingWritingMessage)), "Writing a {count} byte message to connection.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _longPollingDisconnected = private static readonly Action<ILogger, Exception> _longPollingDisconnected =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Debug, new EventId(4, nameof(LongPollingDisconnected)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Client disconnected from Long Polling endpoint for connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(4, nameof(LongPollingDisconnected)), "Client disconnected from Long Polling endpoint for connection.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _longPollingTerminated = private static readonly Action<ILogger, Exception> _longPollingTerminated =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Error, new EventId(5, nameof(LongPollingTerminated)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Long Polling transport was terminated due to an error on connection."); LoggerMessage.Define(LogLevel.Error, new EventId(5, nameof(LongPollingTerminated)), "Long Polling transport was terminated due to an error on connection.");
// Category: HttpConnectionDispatcher // Category: HttpConnectionDispatcher
private static readonly Action<ILogger, DateTime, string, Exception> _connectionDisposed = private static readonly Action<ILogger, string, Exception> _connectionDisposed =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(1, nameof(ConnectionDisposed)), "{time}: Connection Id {connectionId} was disposed."); LoggerMessage.Define<string>(LogLevel.Debug, new EventId(1, nameof(ConnectionDisposed)), "Connection Id {connectionId} was disposed.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _connectionAlreadyActive = private static readonly Action<ILogger, string, string, Exception> _connectionAlreadyActive =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Debug, new EventId(2, nameof(ConnectionAlreadyActive)), "{time}: Connection Id {connectionId} is already active via {requestId}."); LoggerMessage.Define<string, string>(LogLevel.Debug, new EventId(2, nameof(ConnectionAlreadyActive)), "Connection Id {connectionId} is already active via {requestId}.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _pollCanceled = private static readonly Action<ILogger, string, string, Exception> _pollCanceled =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Debug, new EventId(3, nameof(PollCanceled)), "{time}: Previous poll canceled for {connectionId} on {requestId}."); LoggerMessage.Define<string, string>(LogLevel.Debug, new EventId(3, nameof(PollCanceled)), "Previous poll canceled for {connectionId} on {requestId}.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _establishedConnection = private static readonly Action<ILogger, Exception> _establishedConnection =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Debug, new EventId(4, nameof(EstablishedConnection)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Establishing new connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(4, nameof(EstablishedConnection)), "Establishing new connection.");
private static readonly Action<ILogger, DateTime, string, string, Exception> _resumingConnection = private static readonly Action<ILogger, Exception> _resumingConnection =
LoggerMessage.Define<DateTime, string, string>(LogLevel.Debug, new EventId(5, nameof(ResumingConnection)), "{time}: Connection Id {connectionId}, Request Id {requestId}: Resuming existing connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(5, nameof(ResumingConnection)), "Resuming existing connection.");
private static readonly Action<ILogger, DateTime, string, int, Exception> _receivedBytes = private static readonly Action<ILogger, int, Exception> _receivedBytes =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(6, nameof(ReceivedBytes)), "{time}: Connection Id {connectionId}: Received {count} bytes."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(6, nameof(ReceivedBytes)), "Received {count} bytes.");
private static readonly Action<ILogger, DateTime, string, TransportType, Exception> _transportNotSupported = private static readonly Action<ILogger, TransportType, Exception> _transportNotSupported =
LoggerMessage.Define<DateTime, string, TransportType>(LogLevel.Debug, new EventId(7, nameof(TransportNotSupported)), "{time}: Connection Id {connectionId}: {transportType} transport not supported by this endpoint type."); LoggerMessage.Define<TransportType>(LogLevel.Debug, new EventId(7, nameof(TransportNotSupported)), "{transportType} transport not supported by this endpoint type.");
private static readonly Action<ILogger, DateTime, string, TransportType, TransportType, Exception> _cannotChangeTransport = private static readonly Action<ILogger, TransportType, TransportType, Exception> _cannotChangeTransport =
LoggerMessage.Define<DateTime, string, TransportType, TransportType>(LogLevel.Debug, new EventId(8, nameof(CannotChangeTransport)), "{time}: Connection Id {connectionId}: Cannot change transports mid-connection; currently using {transportType}, requesting {requestedTransport}."); LoggerMessage.Define<TransportType, TransportType>(LogLevel.Debug, new EventId(8, nameof(CannotChangeTransport)), "Cannot change transports mid-connection; currently using {transportType}, requesting {requestedTransport}.");
private static readonly Action<ILogger, DateTime, string, Exception> _postNotallowedForWebsockets = private static readonly Action<ILogger, Exception> _postNotallowedForWebsockets =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(9, nameof(PostNotAllowedForWebSockets)), "{time}: Connection Id {connectionId}: POST requests are not allowed for websocket connections."); LoggerMessage.Define(LogLevel.Debug, new EventId(9, nameof(PostNotAllowedForWebSockets)), "POST requests are not allowed for websocket connections.");
private static readonly Action<ILogger, DateTime, string, Exception> _negotiationRequest = private static readonly Action<ILogger, Exception> _negotiationRequest =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(10, nameof(NegotiationRequest)), "{time}: Connection Id {connectionId}: Sending negotiation response."); LoggerMessage.Define(LogLevel.Debug, new EventId(10, nameof(NegotiationRequest)), "Sending negotiation response.");
// Category: WebSocketsTransport // Category: WebSocketsTransport
private static readonly Action<ILogger, DateTime, string, Exception> _socketOpened = private static readonly Action<ILogger, Exception> _socketOpened =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(1, nameof(SocketOpened)), "{time}: Connection Id {connectionId}: Socket opened."); LoggerMessage.Define(LogLevel.Information, new EventId(1, nameof(SocketOpened)), "Socket opened.");
private static readonly Action<ILogger, DateTime, string, Exception> _socketClosed = private static readonly Action<ILogger, Exception> _socketClosed =
LoggerMessage.Define<DateTime, string>(LogLevel.Information, new EventId(2, nameof(SocketClosed)), "{time}: Connection Id {connectionId}: Socket closed."); LoggerMessage.Define(LogLevel.Information, new EventId(2, nameof(SocketClosed)), "Socket closed.");
private static readonly Action<ILogger, DateTime, string, WebSocketCloseStatus?, string, Exception> _clientClosed = private static readonly Action<ILogger, WebSocketCloseStatus?, string, Exception> _clientClosed =
LoggerMessage.Define<DateTime, string, WebSocketCloseStatus?, string>(LogLevel.Debug, new EventId(3, nameof(ClientClosed)), "{time}: Connection Id {connectionId}: Client closed connection with status code '{status}' ({description}). Signaling end-of-input to application."); LoggerMessage.Define<WebSocketCloseStatus?, string>(LogLevel.Debug, new EventId(3, nameof(ClientClosed)), "Client closed connection with status code '{status}' ({description}). Signaling end-of-input to application.");
private static readonly Action<ILogger, DateTime, string, Exception> _waitingForSend = private static readonly Action<ILogger, Exception> _waitingForSend =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(4, nameof(WaitingForSend)), "{time}: Connection Id {connectionId}: Waiting for the application to finish sending data."); LoggerMessage.Define(LogLevel.Debug, new EventId(4, nameof(WaitingForSend)), "Waiting for the application to finish sending data.");
private static readonly Action<ILogger, DateTime, string, Exception> _failedSending = private static readonly Action<ILogger, Exception> _failedSending =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(5, nameof(FailedSending)), "{time}: Connection Id {connectionId}: Application failed during sending. Sending InternalServerError close frame."); LoggerMessage.Define(LogLevel.Debug, new EventId(5, nameof(FailedSending)), "Application failed during sending. Sending InternalServerError close frame.");
private static readonly Action<ILogger, DateTime, string, Exception> _finishedSending = private static readonly Action<ILogger, Exception> _finishedSending =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(6, nameof(FinishedSending)), "{time}: Connection Id {connectionId}: Application finished sending. Sending close frame."); LoggerMessage.Define(LogLevel.Debug, new EventId(6, nameof(FinishedSending)), "Application finished sending. Sending close frame.");
private static readonly Action<ILogger, DateTime, string, Exception> _waitingForClose = private static readonly Action<ILogger, Exception> _waitingForClose =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(7, nameof(WaitingForClose)), "{time}: Connection Id {connectionId}: Waiting for the client to close the socket."); LoggerMessage.Define(LogLevel.Debug, new EventId(7, nameof(WaitingForClose)), "Waiting for the client to close the socket.");
private static readonly Action<ILogger, DateTime, string, Exception> _closeTimedOut = private static readonly Action<ILogger, Exception> _closeTimedOut =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(8, nameof(CloseTimedOut)), "{time}: Connection Id {connectionId}: Timed out waiting for client to send the close frame, aborting the connection."); LoggerMessage.Define(LogLevel.Debug, new EventId(8, nameof(CloseTimedOut)), "Timed out waiting for client to send the close frame, aborting the connection.");
private static readonly Action<ILogger, DateTime, string, WebSocketMessageType, int, bool, Exception> _messageReceived = private static readonly Action<ILogger, WebSocketMessageType, int, bool, Exception> _messageReceived =
LoggerMessage.Define<DateTime, string, WebSocketMessageType, int, bool>(LogLevel.Debug, new EventId(9, nameof(MessageReceived)), "{time}: Connection Id {connectionId}: Message received. Type: {messageType}, size: {size}, EndOfMessage: {endOfMessage}."); LoggerMessage.Define<WebSocketMessageType, int, bool>(LogLevel.Debug, new EventId(9, nameof(MessageReceived)), "Message received. Type: {messageType}, size: {size}, EndOfMessage: {endOfMessage}.");
private static readonly Action<ILogger, DateTime, string, int, Exception> _messageToApplication = private static readonly Action<ILogger, int, Exception> _messageToApplication =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(10, nameof(MessageToApplication)), "{time}: Connection Id {connectionId}: Passing message to application. Payload size: {size}."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(10, nameof(MessageToApplication)), "Passing message to application. Payload size: {size}.");
private static readonly Action<ILogger, DateTime, string, int, Exception> _sendPayload = private static readonly Action<ILogger, int, Exception> _sendPayload =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(11, nameof(SendPayload)), "{time}: Connection Id {connectionId}: Sending payload: {size} bytes."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(11, nameof(SendPayload)), "Sending payload: {size} bytes.");
private static readonly Action<ILogger, DateTime, string, Exception> _errorWritingFrame = private static readonly Action<ILogger, Exception> _errorWritingFrame =
LoggerMessage.Define<DateTime, string>(LogLevel.Error, new EventId(12, nameof(ErrorWritingFrame)), "{time}: Connection Id {connectionId}: Error writing frame."); LoggerMessage.Define(LogLevel.Error, new EventId(12, nameof(ErrorWritingFrame)), "Error writing frame.");
private static readonly Action<ILogger, DateTime, string, Exception> _sendFailed = private static readonly Action<ILogger, Exception> _sendFailed =
LoggerMessage.Define<DateTime, string>(LogLevel.Trace, new EventId(13, nameof(SendFailed)), "{time}: Connection Id {connectionId}: Socket failed to send."); LoggerMessage.Define(LogLevel.Trace, new EventId(13, nameof(SendFailed)), "Socket failed to send.");
// Category: ServerSentEventsTransport // Category: ServerSentEventsTransport
private static readonly Action<ILogger, DateTime, string, int, Exception> _sseWritingMessage = private static readonly Action<ILogger, int, Exception> _sseWritingMessage =
LoggerMessage.Define<DateTime, string, int>(LogLevel.Debug, new EventId(1, nameof(SSEWritingMessage)), "{time}: Connection Id {connectionId}: Writing a {count} byte message."); LoggerMessage.Define<int>(LogLevel.Debug, new EventId(1, nameof(SSEWritingMessage)), "Writing a {count} byte message.");
public static void LongPolling204(this ILogger logger, string connectionId, string requestId) public static void LongPolling204(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _longPolling204(logger, null);
{
_longPolling204(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void PollTimedOut(this ILogger logger, string connectionId, string requestId) public static void PollTimedOut(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _pollTimedOut(logger, null);
{
_pollTimedOut(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void LongPollingWritingMessage(this ILogger logger, string connectionId, string requestId, int count) public static void LongPollingWritingMessage(this ILogger logger, int count)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _longPollingWritingMessage(logger, count, null);
{
_longPollingWritingMessage(logger, DateTime.Now, connectionId, requestId, count, null);
}
} }
public static void LongPollingDisconnected(this ILogger logger, string connectionId, string requestId) public static void LongPollingDisconnected(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _longPollingDisconnected(logger, null);
{
_longPollingDisconnected(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void LongPollingTerminated(this ILogger logger, string connectionId, string requestId, Exception ex) public static void LongPollingTerminated(this ILogger logger, Exception ex)
{ {
if (logger.IsEnabled(LogLevel.Error)) _longPollingTerminated(logger, ex);
{
_longPollingTerminated(logger, DateTime.Now, connectionId, requestId, ex);
}
} }
public static void ConnectionDisposed(this ILogger logger, string connectionId) public static void ConnectionDisposed(this ILogger logger, string connectionId)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _connectionDisposed(logger, connectionId, null);
{
_connectionDisposed(logger, DateTime.Now, connectionId, null);
}
} }
public static void ConnectionAlreadyActive(this ILogger logger, string connectionId, string requestId) public static void ConnectionAlreadyActive(this ILogger logger, string connectionId, string requestId)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _connectionAlreadyActive(logger, connectionId, requestId, null);
{
_connectionAlreadyActive(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void PollCanceled(this ILogger logger, string connectionId, string requestId) public static void PollCanceled(this ILogger logger, string connectionId, string requestId)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _pollCanceled(logger, connectionId, requestId, null);
{
_pollCanceled(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void EstablishedConnection(this ILogger logger, string connectionId, string requestId) public static void EstablishedConnection(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _establishedConnection(logger, null);
{
_establishedConnection(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void ResumingConnection(this ILogger logger, string connectionId, string requestId) public static void ResumingConnection(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _resumingConnection(logger, null);
{
_resumingConnection(logger, DateTime.Now, connectionId, requestId, null);
}
} }
public static void ReceivedBytes(this ILogger logger, string connectionId, int count) public static void ReceivedBytes(this ILogger logger, int count)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _receivedBytes(logger, count, null);
{
_receivedBytes(logger, DateTime.Now, connectionId, count, null);
}
} }
public static void TransportNotSupported(this ILogger logger, string connectionId, TransportType transport) public static void TransportNotSupported(this ILogger logger, TransportType transport)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _transportNotSupported(logger, transport, null);
{
_transportNotSupported(logger, DateTime.Now, connectionId, transport, null);
}
} }
public static void CannotChangeTransport(this ILogger logger, string connectionId, TransportType transport, TransportType requestTransport) public static void CannotChangeTransport(this ILogger logger, TransportType transport, TransportType requestTransport)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _cannotChangeTransport(logger, transport, requestTransport, null);
{
_cannotChangeTransport(logger, DateTime.Now, connectionId, transport, requestTransport, null);
}
} }
public static void PostNotAllowedForWebSockets(this ILogger logger, string connectionId) public static void PostNotAllowedForWebSockets(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _postNotallowedForWebsockets(logger, null);
{
_postNotallowedForWebsockets(logger, DateTime.Now, connectionId, null);
}
} }
public static void NegotiationRequest(this ILogger logger, string connectionId) public static void NegotiationRequest(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _negotiationRequest(logger, null);
{
_negotiationRequest(logger, DateTime.Now, connectionId, null);
}
} }
public static void SocketOpened(this ILogger logger, string connectionId) public static void SocketOpened(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _socketOpened(logger, null);
{
_socketOpened(logger, DateTime.Now, connectionId, null);
}
} }
public static void SocketClosed(this ILogger logger, string connectionId) public static void SocketClosed(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Information)) _socketClosed(logger, null);
{
_socketClosed(logger, DateTime.Now, connectionId, null);
}
} }
public static void ClientClosed(this ILogger logger, string connectionId, WebSocketCloseStatus? closeStatus, string closeDescription) public static void ClientClosed(this ILogger logger, WebSocketCloseStatus? closeStatus, string closeDescription)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _clientClosed(logger, closeStatus, closeDescription, null);
{
_clientClosed(logger, DateTime.Now, connectionId, closeStatus, closeDescription, null);
}
} }
public static void WaitingForSend(this ILogger logger, string connectionId) public static void WaitingForSend(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _waitingForSend(logger, null);
{
_waitingForSend(logger, DateTime.Now, connectionId, null);
}
} }
public static void FailedSending(this ILogger logger, string connectionId) public static void FailedSending(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _failedSending(logger, null);
{
_failedSending(logger, DateTime.Now, connectionId, null);
}
} }
public static void FinishedSending(this ILogger logger, string connectionId) public static void FinishedSending(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _finishedSending(logger, null);
{
_finishedSending(logger, DateTime.Now, connectionId, null);
}
} }
public static void WaitingForClose(this ILogger logger, string connectionId) public static void WaitingForClose(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _waitingForClose(logger, null);
{
_waitingForClose(logger, DateTime.Now, connectionId, null);
}
} }
public static void CloseTimedOut(this ILogger logger, string connectionId) public static void CloseTimedOut(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _closeTimedOut(logger, null);
{
_closeTimedOut(logger, DateTime.Now, connectionId, null);
}
} }
public static void MessageReceived(this ILogger logger, string connectionId, WebSocketMessageType type, int size, bool endOfMessage) public static void MessageReceived(this ILogger logger, WebSocketMessageType type, int size, bool endOfMessage)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _messageReceived(logger, type, size, endOfMessage, null);
{
_messageReceived(logger, DateTime.Now, connectionId, type, size, endOfMessage, null);
}
} }
public static void MessageToApplication(this ILogger logger, string connectionId, int size) public static void MessageToApplication(this ILogger logger, int size)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _messageToApplication(logger, size, null);
{
_messageToApplication(logger, DateTime.Now, connectionId, size, null);
}
} }
public static void SendPayload(this ILogger logger, string connectionId, int size) public static void SendPayload(this ILogger logger, int size)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sendPayload(logger, size, null);
{
_sendPayload(logger, DateTime.Now, connectionId, size, null);
}
} }
public static void ErrorWritingFrame(this ILogger logger, string connectionId, Exception ex) public static void ErrorWritingFrame(this ILogger logger, Exception ex)
{ {
if (logger.IsEnabled(LogLevel.Error)) _errorWritingFrame(logger, ex);
{
_errorWritingFrame(logger, DateTime.Now, connectionId, ex);
}
} }
public static void SendFailed(this ILogger logger, string connectionId, Exception ex) public static void SendFailed(this ILogger logger, Exception ex)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _sendFailed(logger, ex);
{
_sendFailed(logger, DateTime.Now, connectionId, ex);
}
} }
public static void SSEWritingMessage(this ILogger logger, string connectionId, int count) public static void SSEWritingMessage(this ILogger logger, int count)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _sseWritingMessage(logger, count, null);
{
_sseWritingMessage(logger, DateTime.Now, connectionId, count, null);
}
} }
} }
} }

View File

@ -34,7 +34,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
if (!await _application.WaitToReadAsync(token)) if (!await _application.WaitToReadAsync(token))
{ {
await _application.Completion; await _application.Completion;
_logger.LongPolling204(_connectionId, context.TraceIdentifier); _logger.LongPolling204();
context.Response.ContentType = "text/plain"; context.Response.ContentType = "text/plain";
context.Response.StatusCode = StatusCodes.Status204NoContent; context.Response.StatusCode = StatusCodes.Status204NoContent;
return; return;
@ -49,7 +49,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
contentLength += buffer.Length; contentLength += buffer.Length;
buffers.Add(buffer); buffers.Add(buffer);
_logger.LongPollingWritingMessage(_connectionId, context.TraceIdentifier, buffer.Length); _logger.LongPollingWritingMessage(buffer.Length);
} }
context.Response.ContentLength = contentLength; context.Response.ContentLength = contentLength;
@ -72,12 +72,12 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
{ {
// Don't count this as cancellation, this is normal as the poll can end due to the browser closing. // Don't count this as cancellation, this is normal as the poll can end due to the browser closing.
// The background thread will eventually dispose this connection if it's inactive // The background thread will eventually dispose this connection if it's inactive
_logger.LongPollingDisconnected(_connectionId, context.TraceIdentifier); _logger.LongPollingDisconnected();
} }
// Case 2 // Case 2
else if (_timeoutToken.IsCancellationRequested) else if (_timeoutToken.IsCancellationRequested)
{ {
_logger.PollTimedOut(_connectionId, context.TraceIdentifier); _logger.PollTimedOut();
context.Response.ContentLength = 0; context.Response.ContentLength = 0;
context.Response.ContentType = "text/plain"; context.Response.ContentType = "text/plain";
@ -86,14 +86,14 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
else else
{ {
// Case 3 // Case 3
_logger.LongPolling204(_connectionId, context.TraceIdentifier); _logger.LongPolling204();
context.Response.ContentType = "text/plain"; context.Response.ContentType = "text/plain";
context.Response.StatusCode = StatusCodes.Status204NoContent; context.Response.StatusCode = StatusCodes.Status204NoContent;
} }
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.LongPollingTerminated(_connectionId, context.TraceIdentifier, ex); _logger.LongPollingTerminated(ex);
throw; throw;
} }
} }

View File

@ -49,7 +49,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
var ms = new MemoryStream(); var ms = new MemoryStream();
while (_application.TryRead(out var buffer)) while (_application.TryRead(out var buffer))
{ {
_logger.SSEWritingMessage(_connectionId, buffer.Length); _logger.SSEWritingMessage(buffer.Length);
ServerSentEventsMessageFormatter.WriteMessage(buffer, ms); ServerSentEventsMessageFormatter.WriteMessage(buffer, ms);
} }

View File

@ -49,7 +49,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
using (var ws = await context.WebSockets.AcceptWebSocketAsync(_options.SubProtocol)) using (var ws = await context.WebSockets.AcceptWebSocketAsync(_options.SubProtocol))
{ {
_logger.SocketOpened(_connection.ConnectionId); _logger.SocketOpened();
try try
{ {
@ -57,7 +57,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
} }
finally finally
{ {
_logger.SocketClosed(_connection.ConnectionId); _logger.SocketClosed();
} }
} }
} }
@ -78,12 +78,12 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
if (trigger == receiving) if (trigger == receiving)
{ {
task = sending; task = sending;
_logger.WaitingForSend(_connection.ConnectionId); _logger.WaitingForSend();
} }
else else
{ {
task = receiving; task = receiving;
_logger.WaitingForClose(_connection.ConnectionId); _logger.WaitingForClose();
} }
// We're done writing // We're done writing
@ -95,7 +95,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
if (resultTask != task) if (resultTask != task)
{ {
_logger.CloseTimedOut(_connection.ConnectionId); _logger.CloseTimedOut();
socket.Abort(); socket.Abort();
} }
else else
@ -129,7 +129,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
return receiveResult; return receiveResult;
} }
_logger.MessageReceived(_connection.ConnectionId, receiveResult.MessageType, receiveResult.Count, receiveResult.EndOfMessage); _logger.MessageReceived(receiveResult.MessageType, receiveResult.Count, receiveResult.EndOfMessage);
var truncBuffer = new ArraySegment<byte>(buffer.Array, 0, receiveResult.Count); var truncBuffer = new ArraySegment<byte>(buffer.Array, 0, receiveResult.Count);
incomingMessage.Add(truncBuffer); incomingMessage.Add(truncBuffer);
@ -159,7 +159,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
Buffer.BlockCopy(incomingMessage[0].Array, incomingMessage[0].Offset, messageBuffer, 0, incomingMessage[0].Count); Buffer.BlockCopy(incomingMessage[0].Array, incomingMessage[0].Offset, messageBuffer, 0, incomingMessage[0].Count);
} }
_logger.MessageToApplication(_connection.ConnectionId, messageBuffer.Length); _logger.MessageToApplication(messageBuffer.Length);
while (await _application.Writer.WaitToWriteAsync()) while (await _application.Writer.WaitToWriteAsync())
{ {
if (_application.Writer.TryWrite(messageBuffer)) if (_application.Writer.TryWrite(messageBuffer))
@ -182,7 +182,7 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
{ {
try try
{ {
_logger.SendPayload(_connection.ConnectionId, buffer.Length); _logger.SendPayload(buffer.Length);
var webSocketMessageType = (_connection.TransferMode == TransferMode.Binary var webSocketMessageType = (_connection.TransferMode == TransferMode.Binary
? WebSocketMessageType.Binary ? WebSocketMessageType.Binary
@ -196,12 +196,12 @@ namespace Microsoft.AspNetCore.Sockets.Internal.Transports
catch (WebSocketException socketException) when (!WebSocketCanSend(ws)) catch (WebSocketException socketException) when (!WebSocketCanSend(ws))
{ {
// this can happen when we send the CloseFrame to the client and try to write afterwards // this can happen when we send the CloseFrame to the client and try to write afterwards
_logger.SendFailed(_connection.ConnectionId, socketException); _logger.SendFailed(socketException);
break; break;
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.ErrorWritingFrame(_connection.ConnectionId, ex); _logger.ErrorWritingFrame(ex);
break; break;
} }
} }

View File

@ -9,81 +9,60 @@ namespace Microsoft.AspNetCore.Sockets.Internal
internal static class SocketLoggerExtensions internal static class SocketLoggerExtensions
{ {
// Category: ConnectionManager // Category: ConnectionManager
private static readonly Action<ILogger, DateTime, string, Exception> _createdNewConnection = private static readonly Action<ILogger, string, Exception> _createdNewConnection =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(1, nameof(CreatedNewConnection)), "{time}: ConnectionId {connectionId}: New connection created."); LoggerMessage.Define<string>(LogLevel.Debug, new EventId(1, nameof(CreatedNewConnection)), "New connection {connectionId} created.");
private static readonly Action<ILogger, DateTime, string, Exception> _removedConnection = private static readonly Action<ILogger, string, Exception> _removedConnection =
LoggerMessage.Define<DateTime, string>(LogLevel.Debug, new EventId(2, nameof(RemovedConnection)), "{time}: ConnectionId {connectionId}: Removing connection from the list of connections."); LoggerMessage.Define<string>(LogLevel.Debug, new EventId(2, nameof(RemovedConnection)), "Removing connection {connectionId} from the list of connections.");
private static readonly Action<ILogger, DateTime, string, Exception> _failedDispose = private static readonly Action<ILogger, string, Exception> _failedDispose =
LoggerMessage.Define<DateTime, string>(LogLevel.Error, new EventId(3, nameof(FailedDispose)), "{time}: ConnectionId {connectionId}: Failed disposing connection."); LoggerMessage.Define<string>(LogLevel.Error, new EventId(3, nameof(FailedDispose)), "Failed disposing connection {connectionId}.");
private static readonly Action<ILogger, DateTime, string, Exception> _connectionReset = private static readonly Action<ILogger, string, Exception> _connectionReset =
LoggerMessage.Define<DateTime, string>(LogLevel.Trace, new EventId(4, nameof(ConnectionReset)), "{time}: ConnectionId {connectionId}: Connection was reset."); LoggerMessage.Define<string>(LogLevel.Trace, new EventId(4, nameof(ConnectionReset)), "Connection {connectionId} was reset.");
private static readonly Action<ILogger, DateTime, string, Exception> _connectionTimedOut = private static readonly Action<ILogger, string, Exception> _connectionTimedOut =
LoggerMessage.Define<DateTime, string>(LogLevel.Trace, new EventId(5, nameof(ConnectionTimedOut)), "{time}: ConnectionId {connectionId}: Connection timed out."); LoggerMessage.Define<string>(LogLevel.Trace, new EventId(5, nameof(ConnectionTimedOut)), "Connection {connectionId} timed out.");
private static readonly Action<ILogger, DateTime, Exception> _scanningConnections = private static readonly Action<ILogger, Exception> _scanningConnections =
LoggerMessage.Define<DateTime>(LogLevel.Trace, new EventId(6, nameof(ScanningConnections)), "{time}: Scanning connections."); LoggerMessage.Define(LogLevel.Trace, new EventId(6, nameof(ScanningConnections)), "Scanning connections.");
private static readonly Action<ILogger, DateTime, TimeSpan, Exception> _scannedConnections = private static readonly Action<ILogger, TimeSpan, Exception> _scannedConnections =
LoggerMessage.Define<DateTime, TimeSpan>(LogLevel.Trace, new EventId(7, nameof(ScannedConnections)), "{time}: Scanned connections in {duration}."); LoggerMessage.Define<TimeSpan>(LogLevel.Trace, new EventId(7, nameof(ScannedConnections)), "Scanned connections in {duration}.");
public static void CreatedNewConnection(this ILogger logger, string connectionId) public static void CreatedNewConnection(this ILogger logger, string connectionId)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _createdNewConnection(logger, connectionId, null);
{
_createdNewConnection(logger, DateTime.Now, connectionId, null);
}
} }
public static void RemovedConnection(this ILogger logger, string connectionId) public static void RemovedConnection(this ILogger logger, string connectionId)
{ {
if (logger.IsEnabled(LogLevel.Debug)) _removedConnection(logger, connectionId, null);
{
_removedConnection(logger, DateTime.Now, connectionId, null);
}
} }
public static void FailedDispose(this ILogger logger, string connectionId, Exception exception) public static void FailedDispose(this ILogger logger, string connectionId, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Error)) _failedDispose(logger, connectionId, exception);
{
_failedDispose(logger, DateTime.Now, connectionId, exception);
}
} }
public static void ConnectionTimedOut(this ILogger logger, string connectionId) public static void ConnectionTimedOut(this ILogger logger, string connectionId)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _connectionTimedOut(logger, connectionId, null);
{
_connectionTimedOut(logger, DateTime.Now, connectionId, null);
}
} }
public static void ConnectionReset(this ILogger logger, string connectionId, Exception exception) public static void ConnectionReset(this ILogger logger, string connectionId, Exception exception)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _connectionReset(logger, connectionId, exception);
{
_connectionReset(logger, DateTime.Now, connectionId, exception);
}
} }
public static void ScanningConnections(this ILogger logger) public static void ScanningConnections(this ILogger logger)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _scanningConnections(logger, null);
{
_scanningConnections(logger, DateTime.Now, null);
}
} }
public static void ScannedConnections(this ILogger logger, TimeSpan duration) public static void ScannedConnections(this ILogger logger, TimeSpan duration)
{ {
if (logger.IsEnabled(LogLevel.Trace)) _scannedConnections(logger, duration, null);
{
_scannedConnections(logger, DateTime.Now, duration, null);
}
} }
} }
} }

View File

@ -44,7 +44,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
transportActiveTask = longPollingTransport.Running; transportActiveTask = longPollingTransport.Running;
@ -80,7 +80,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = ChannelConnection.Create(connectionToTransport, transportToConnection); var channelConnection = ChannelConnection.Create(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
await longPollingTransport.Running.OrTimeout(); await longPollingTransport.Running.OrTimeout();
Assert.True(transportToConnection.Reader.Completion.IsCompleted); Assert.True(transportToConnection.Reader.Completion.IsCompleted);
@ -133,7 +133,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
var data = await transportToConnection.Reader.ReadAllAsync().OrTimeout(); var data = await transportToConnection.Reader.ReadAllAsync().OrTimeout();
await longPollingTransport.Running.OrTimeout(); await longPollingTransport.Running.OrTimeout();
@ -169,7 +169,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
var exception = var exception =
await Assert.ThrowsAsync<HttpRequestException>(async () => await transportToConnection.Reader.Completion.OrTimeout()); await Assert.ThrowsAsync<HttpRequestException>(async () => await transportToConnection.Reader.Completion.OrTimeout());
@ -205,7 +205,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
await connectionToTransport.Writer.WriteAsync(new SendMessage()); await connectionToTransport.Writer.WriteAsync(new SendMessage());
@ -246,7 +246,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
connectionToTransport.Writer.Complete(); connectionToTransport.Writer.Complete();
@ -297,7 +297,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
// Start the transport // Start the transport
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
// Wait for the transport to finish // Wait for the transport to finish
await longPollingTransport.Running.OrTimeout(); await longPollingTransport.Running.OrTimeout();
@ -362,7 +362,7 @@ namespace Microsoft.AspNetCore.Client.Tests
await connectionToTransport.Writer.WriteAsync(new SendMessage(Encoding.UTF8.GetBytes("World"), tcs2)).OrTimeout(); await connectionToTransport.Writer.WriteAsync(new SendMessage(Encoding.UTF8.GetBytes("World"), tcs2)).OrTimeout();
// Start the transport // Start the transport
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
connectionToTransport.Writer.Complete(); connectionToTransport.Writer.Complete();
@ -405,7 +405,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
Assert.Null(longPollingTransport.Mode); Assert.Null(longPollingTransport.Mode);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, transferMode, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, transferMode, connection: new TestConnection());
Assert.Equal(transferMode, longPollingTransport.Mode); Assert.Equal(transferMode, longPollingTransport.Mode);
} }
finally finally
@ -431,7 +431,7 @@ namespace Microsoft.AspNetCore.Client.Tests
{ {
var longPollingTransport = new LongPollingTransport(httpClient); var longPollingTransport = new LongPollingTransport(httpClient);
var exception = await Assert.ThrowsAsync<ArgumentException>(() => var exception = await Assert.ThrowsAsync<ArgumentException>(() =>
longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), null, TransferMode.Text | TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection())); longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), null, TransferMode.Text | TransferMode.Binary, connection: new TestConnection()));
Assert.Contains("Invalid transfer mode.", exception.Message); Assert.Contains("Invalid transfer mode.", exception.Message);
Assert.Equal("requestedTransferMode", exception.ParamName); Assert.Equal("requestedTransferMode", exception.ParamName);
@ -469,7 +469,7 @@ namespace Microsoft.AspNetCore.Client.Tests
var connectionToTransport = Channel.CreateUnbounded<SendMessage>(); var connectionToTransport = Channel.CreateUnbounded<SendMessage>();
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connectionId: string.Empty, connection: new TestConnection()); await longPollingTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Binary, connection: new TestConnection());
var completedTask = await Task.WhenAny(completionTcs.Task, longPollingTransport.Running).OrTimeout(); var completedTask = await Task.WhenAny(completionTcs.Task, longPollingTransport.Running).OrTimeout();
Assert.Equal(completionTcs.Task, completedTask); Assert.Equal(completionTcs.Task, completedTask);

View File

@ -55,7 +55,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
await eventStreamTcs.Task.OrTimeout(); await eventStreamTcs.Task.OrTimeout();
await sseTransport.StopAsync().OrTimeout(); await sseTransport.StopAsync().OrTimeout();
@ -106,7 +106,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
transportActiveTask = sseTransport.Running; transportActiveTask = sseTransport.Running;
Assert.False(transportActiveTask.IsCompleted); Assert.False(transportActiveTask.IsCompleted);
@ -154,7 +154,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
var exception = await Assert.ThrowsAsync<FormatException>(() => sseTransport.Running.OrTimeout()); var exception = await Assert.ThrowsAsync<FormatException>(() => sseTransport.Running.OrTimeout());
Assert.Equal("Incomplete message.", exception.Message); Assert.Equal("Incomplete message.", exception.Message);
@ -200,7 +200,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
await eventStreamTcs.Task; await eventStreamTcs.Task;
var sendTcs = new TaskCompletionSource<object>(); var sendTcs = new TaskCompletionSource<object>();
@ -247,7 +247,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
await eventStreamTcs.Task.OrTimeout(); await eventStreamTcs.Task.OrTimeout();
connectionToTransport.Writer.TryComplete(null); connectionToTransport.Writer.TryComplete(null);
@ -276,7 +276,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
await sseTransport.StartAsync( await sseTransport.StartAsync(
new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text, connection: Mock.Of<IConnection>()).OrTimeout();
var message = await transportToConnection.Reader.ReadAsync().AsTask().OrTimeout(); var message = await transportToConnection.Reader.ReadAsync().AsTask().OrTimeout();
Assert.Equal("3:abc", Encoding.ASCII.GetString(message)); Assert.Equal("3:abc", Encoding.ASCII.GetString(message));
@ -306,7 +306,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
Assert.Null(sseTransport.Mode); Assert.Null(sseTransport.Mode);
await sseTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, transferMode, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); await sseTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, transferMode, connection: Mock.Of<IConnection>()).OrTimeout();
Assert.Equal(TransferMode.Text, sseTransport.Mode); Assert.Equal(TransferMode.Text, sseTransport.Mode);
await sseTransport.StopAsync().OrTimeout(); await sseTransport.StopAsync().OrTimeout();
} }
@ -331,7 +331,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
var transportToConnection = Channel.CreateUnbounded<byte[]>(); var transportToConnection = Channel.CreateUnbounded<byte[]>();
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
var exception = await Assert.ThrowsAsync<ArgumentException>(() => var exception = await Assert.ThrowsAsync<ArgumentException>(() =>
sseTransport.StartAsync(new Uri("http://fakeuri.org"), null, TransferMode.Text | TransferMode.Binary, connectionId: string.Empty, connection: Mock.Of<IConnection>())); sseTransport.StartAsync(new Uri("http://fakeuri.org"), null, TransferMode.Text | TransferMode.Binary, connection: Mock.Of<IConnection>()));
Assert.Contains("Invalid transfer mode.", exception.Message); Assert.Contains("Invalid transfer mode.", exception.Message);
Assert.Equal("requestedTransferMode", exception.ParamName); Assert.Equal("requestedTransferMode", exception.ParamName);

View File

@ -21,7 +21,7 @@ namespace Microsoft.AspNetCore.SignalR.Client.Tests
Mode = transferMode; Mode = transferMode;
} }
public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, string connectionId, IConnection connection) public Task StartAsync(Uri url, Channel<byte[], SendMessage> application, TransferMode requestedTransferMode, IConnection connection)
{ {
Application = application; Application = application;
return _startHandler(); return _startHandler();

View File

@ -42,7 +42,7 @@ namespace Microsoft.AspNetCore.SignalR.Tests
var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory); var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory);
await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection, await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection,
TransferMode.Binary, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); TransferMode.Binary, connection: Mock.Of<IConnection>()).OrTimeout();
await webSocketsTransport.StopAsync().OrTimeout(); await webSocketsTransport.StopAsync().OrTimeout();
await webSocketsTransport.Running.OrTimeout(); await webSocketsTransport.Running.OrTimeout();
} }
@ -60,7 +60,7 @@ namespace Microsoft.AspNetCore.SignalR.Tests
var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory); var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory);
await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection, await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection,
TransferMode.Binary, connectionId: string.Empty, connection: Mock.Of<IConnection>()); TransferMode.Binary, connection: Mock.Of<IConnection>());
connectionToTransport.Writer.TryComplete(); connectionToTransport.Writer.TryComplete();
await webSocketsTransport.Running.OrTimeout(TimeSpan.FromSeconds(10)); await webSocketsTransport.Running.OrTimeout(TimeSpan.FromSeconds(10));
} }
@ -79,7 +79,7 @@ namespace Microsoft.AspNetCore.SignalR.Tests
var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection); var channelConnection = new ChannelConnection<SendMessage, byte[]>(connectionToTransport, transportToConnection);
var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory); var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory);
await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection, transferMode, connectionId: string.Empty, connection: Mock.Of<IConnection>()); await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection, transferMode, connection: Mock.Of<IConnection>());
var sendTcs = new TaskCompletionSource<object>(); var sendTcs = new TaskCompletionSource<object>();
connectionToTransport.Writer.TryWrite(new SendMessage(new byte[] { 0x42 }, sendTcs)); connectionToTransport.Writer.TryWrite(new SendMessage(new byte[] { 0x42 }, sendTcs));
@ -120,7 +120,7 @@ namespace Microsoft.AspNetCore.SignalR.Tests
Assert.Null(webSocketsTransport.Mode); Assert.Null(webSocketsTransport.Mode);
await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection, await webSocketsTransport.StartAsync(new Uri(_serverFixture.WebSocketsUrl + "/echo"), channelConnection,
transferMode, connectionId: string.Empty, connection: Mock.Of<IConnection>()).OrTimeout(); transferMode, connection: Mock.Of<IConnection>()).OrTimeout();
Assert.Equal(transferMode, webSocketsTransport.Mode); Assert.Equal(transferMode, webSocketsTransport.Mode);
await webSocketsTransport.StopAsync().OrTimeout(); await webSocketsTransport.StopAsync().OrTimeout();
@ -140,7 +140,7 @@ namespace Microsoft.AspNetCore.SignalR.Tests
var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory); var webSocketsTransport = new WebSocketsTransport(httpOptions: null, loggerFactory: loggerFactory);
var exception = await Assert.ThrowsAsync<ArgumentException>(() => var exception = await Assert.ThrowsAsync<ArgumentException>(() =>
webSocketsTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text | TransferMode.Binary, connectionId: string.Empty, connection: Mock.Of<IConnection>())); webSocketsTransport.StartAsync(new Uri("http://fakeuri.org"), channelConnection, TransferMode.Text | TransferMode.Binary, connection: Mock.Of<IConnection>()));
Assert.Contains("Invalid transfer mode.", exception.Message); Assert.Contains("Invalid transfer mode.", exception.Message);
Assert.Equal("requestedTransferMode", exception.ParamName); Assert.Equal("requestedTransferMode", exception.ParamName);