1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-15 02:13:01 +00:00

Implement high-performance logging (#193)

* Implement high-performance logging
This commit is contained in:
Jonnern
2024-03-22 16:39:32 +01:00
committed by GitHub
parent 108c8fc183
commit de72fe4fb9
11 changed files with 1293 additions and 112 deletions
@@ -1,4 +1,5 @@
using CryptoExchange.Net.Interfaces;
using CryptoExchange.Net.Logging.Extensions;
using CryptoExchange.Net.Objects;
using CryptoExchange.Net.Objects.Sockets;
using Microsoft.Extensions.Logging;
@@ -189,7 +190,7 @@ namespace CryptoExchange.Net.Sockets
private async Task<bool> ConnectInternalAsync()
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] connecting");
_logger.SocketConnecting(Id);
try
{
using CancellationTokenSource tcs = new(TimeSpan.FromSeconds(10));
@@ -197,11 +198,11 @@ namespace CryptoExchange.Net.Sockets
}
catch (Exception e)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] connection failed: " + e.ToLogString());
_logger.SocketConnectionFailed(Id, e.Message, e);
return false;
}
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] connected to {Uri}");
_logger.SocketConnected(Id, Uri);
return true;
}
@@ -210,13 +211,13 @@ namespace CryptoExchange.Net.Sockets
{
while (!_stopRequested)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] starting processing tasks");
_logger.SocketStartingProcessing(Id);
_processState = ProcessState.Processing;
var sendTask = SendLoopAsync();
var receiveTask = ReceiveLoopAsync();
var timeoutTask = Parameters.Timeout != null && Parameters.Timeout > TimeSpan.FromSeconds(0) ? CheckTimeoutAsync() : Task.CompletedTask;
await Task.WhenAll(sendTask, receiveTask, timeoutTask).ConfigureAwait(false);
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] processing tasks finished");
_logger.SocketFinishedProcessing(Id);
_processState = ProcessState.WaitingForClose;
while (_closeTask == null)
@@ -244,14 +245,14 @@ namespace CryptoExchange.Net.Sockets
while (!_stopRequested)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] attempting to reconnect");
_logger.SocketAttemptReconnect(Id);
var task = GetReconnectionUrl?.Invoke();
if (task != null)
{
var reconnectUri = await task.ConfigureAwait(false);
if (reconnectUri != null && Parameters.Uri != reconnectUri)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] reconnect URI set to {reconnectUri}");
_logger.SocketSetReconnectUri(Id, reconnectUri);
Parameters.Uri = reconnectUri;
}
}
@@ -284,7 +285,7 @@ namespace CryptoExchange.Net.Sockets
return;
var bytes = Parameters.Encoding.GetBytes(data);
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] msg {id} - Adding {bytes.Length} bytes to send buffer");
_logger.SocketAddingBytesToSendBuffer(Id, id, bytes);
_sendBuffer.Enqueue(new SendItem { Id = id, Weight = weight, Bytes = bytes });
_sendEvent.Set();
}
@@ -295,7 +296,7 @@ namespace CryptoExchange.Net.Sockets
if (_processState != ProcessState.Processing && IsOpen)
return;
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] reconnect requested");
_logger.SocketReconnectRequested(Id);
_closeTask = CloseInternalAsync();
await _closeTask.ConfigureAwait(false);
}
@@ -310,18 +311,18 @@ namespace CryptoExchange.Net.Sockets
{
if (_closeTask?.IsCompleted == false)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] CloseAsync() waiting for existing close task");
_logger.SocketCloseAsyncWaitingForExistingCloseTask(Id);
await _closeTask.ConfigureAwait(false);
return;
}
if (!IsOpen)
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] CloseAsync() socket not open");
_logger.SocketCloseAsyncSocketNotOpen(Id);
return;
}
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] closing");
_logger.SocketClosing(Id);
_closeTask = CloseInternalAsync();
}
finally
@@ -333,7 +334,7 @@ namespace CryptoExchange.Net.Sockets
if(_processTask != null)
await _processTask.ConfigureAwait(false);
await (OnClose?.Invoke() ?? Task.CompletedTask).ConfigureAwait(false);
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] closed");
_logger.SocketClosed(Id);
}
/// <summary>
@@ -385,11 +386,11 @@ namespace CryptoExchange.Net.Sockets
if (_disposed)
return;
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] disposing");
_logger.SocketDisposing(Id);
_disposed = true;
_socket.Dispose();
_ctsSource.Dispose();
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] disposed");
_logger.SocketDisposed(Id);
}
/// <summary>
@@ -421,7 +422,7 @@ namespace CryptoExchange.Net.Sockets
if (limitResult.Success)
{
if (limitResult.Data > 0)
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] msg {data.Id} - send delayed {limitResult.Data}ms because of rate limit");
_logger.SocketSendDelayedBecauseOfRateLimit(Id, data.Id, limitResult.Data);
}
}
}
@@ -430,7 +431,7 @@ namespace CryptoExchange.Net.Sockets
{
await _socket.SendAsync(new ArraySegment<byte>(data.Bytes, 0, data.Bytes.Length), WebSocketMessageType.Text, true, _ctsSource.Token).ConfigureAwait(false);
await (OnRequestSent?.Invoke(data.Id) ?? Task.CompletedTask).ConfigureAwait(false);
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] msg {data.Id} - sent {data.Bytes.Length} bytes");
_logger.SocketSentBytes(Id, data.Id, data.Bytes.Length);
}
catch (OperationCanceledException)
{
@@ -453,13 +454,13 @@ namespace CryptoExchange.Net.Sockets
// Because this is running in a separate task and not awaited until the socket gets closed
// any exception here will crash the send processing, but do so silently unless the socket get's stopped.
// Make sure we at least let the owner know there was an error
_logger.Log(LogLevel.Warning, $"[Sckt {Id}] send loop stopped with exception");
_logger.SocketSendLoopStoppedWithException(Id, e.Message, e);
await (OnError?.Invoke(e) ?? Task.CompletedTask).ConfigureAwait(false);
throw;
}
finally
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] send loop finished");
_logger.SocketSendLoopFinished(Id);
}
}
@@ -506,8 +507,8 @@ namespace CryptoExchange.Net.Sockets
if (receiveResult.MessageType == WebSocketMessageType.Close)
{
// Connection closed unexpectedly
_logger.Log(LogLevel.Debug, "[Sckt {Id}] received `Close` message, CloseStatus: {Status}, CloseStatusDescription: {CloseStatusDescription}", Id, receiveResult.CloseStatus, receiveResult.CloseStatusDescription);
// Connection closed unexpectedly
_logger.SocketReceivedCloseMessage(Id, receiveResult.CloseStatus.ToString(), receiveResult.CloseStatusDescription);
if (_closeTask?.IsCompleted != false)
_closeTask = CloseInternalAsync();
break;
@@ -517,7 +518,7 @@ namespace CryptoExchange.Net.Sockets
{
// We received data, but it is not complete, write it to a memory stream for reassembling
multiPartMessage = true;
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] received {receiveResult.Count} bytes in partial message");
_logger.SocketReceivedPartialMessage(Id, receiveResult.Count);
// Write the data to a memory stream to be reassembled later
if (multipartStream == null)
@@ -529,13 +530,13 @@ namespace CryptoExchange.Net.Sockets
if (!multiPartMessage)
{
// Received a complete message and it's not multi part
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] received {receiveResult.Count} bytes in single message");
_logger.SocketReceivedSingleMessage(Id, receiveResult.Count);
ProcessData(receiveResult.MessageType, new ReadOnlyMemory<byte>(buffer.Array, buffer.Offset, receiveResult.Count));
}
else
{
// Received the end of a multipart message, write to memory stream for reassembling
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] received {receiveResult.Count} bytes in partial message");
_logger.SocketReceivedPartialMessage(Id, receiveResult.Count);
multipartStream!.Write(buffer.Array, buffer.Offset, receiveResult.Count);
}
@@ -563,13 +564,13 @@ namespace CryptoExchange.Net.Sockets
// When the connection gets interupted we might not have received a full message
if (receiveResult?.EndOfMessage == true)
{
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] reassembled message of {multipartStream!.Length} bytes");
_logger.SocketReassembledMessage(Id, multipartStream!.Length);
// Get the underlying buffer of the memorystream holding the written data and delimit it (GetBuffer return the full array, not only the written part)
ProcessData(receiveResult.MessageType, new ReadOnlyMemory<byte>(multipartStream.GetBuffer(), 0, (int)multipartStream.Length));
}
else
{
_logger.Log(LogLevel.Trace, $"[Sckt {Id}] discarding incomplete message of {multipartStream!.Length} bytes");
_logger.SocketDiscardIncompleteMessage(Id, multipartStream!.Length);
}
}
}
@@ -579,13 +580,13 @@ namespace CryptoExchange.Net.Sockets
// Because this is running in a separate task and not awaited until the socket gets closed
// any exception here will crash the receive processing, but do so silently unless the socket gets stopped.
// Make sure we at least let the owner know there was an error
_logger.Log(LogLevel.Warning, $"[Sckt {Id}] receive loop stopped with exception");
_logger.SocketReceiveLoopStoppedWithException(Id, e);
await (OnError?.Invoke(e) ?? Task.CompletedTask).ConfigureAwait(false);
throw;
}
finally
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] receive loop finished");
_logger.SocketReceiveLoopFinished(Id);
}
}
@@ -607,7 +608,7 @@ namespace CryptoExchange.Net.Sockets
/// <returns></returns>
protected async Task CheckTimeoutAsync()
{
_logger.Log(LogLevel.Debug, $"[Sckt {Id}] starting task checking for no data received for {Parameters.Timeout}");
_logger.SocketStartingTaskForNoDataReceivedCheck(Id, Parameters.Timeout);
LastActionTime = DateTime.UtcNow;
try
{
@@ -618,7 +619,7 @@ namespace CryptoExchange.Net.Sockets
if (DateTime.UtcNow - LastActionTime > Parameters.Timeout)
{
_logger.Log(LogLevel.Warning, $"[Sckt {Id}] no data received for {Parameters.Timeout}, reconnecting socket");
_logger.SocketNoDataReceiveTimoutReconnect(Id, Parameters.Timeout);
_ = ReconnectAsync().ConfigureAwait(false);
return;
}