diff --git a/EonaCat.LogStack/EonaCat.LogStack.csproj b/EonaCat.LogStack/EonaCat.LogStack.csproj index 9b70ba8..340059d 100644 --- a/EonaCat.LogStack/EonaCat.LogStack.csproj +++ b/EonaCat.LogStack/EonaCat.LogStack.csproj @@ -25,7 +25,7 @@ It features a rich fluent API for routing log events to dozens of destinations f - 0.2.3+{chash:10}.{c:ymd} + 0.2.3+{chash:10}.{c:ymd} true true v[0-9]* diff --git a/EonaCat.LogStack/EonaCatLogger.cs b/EonaCat.LogStack/EonaCatLogger.cs index eb88c52..91b3b71 100644 --- a/EonaCat.LogStack/EonaCatLogger.cs +++ b/EonaCat.LogStack/EonaCatLogger.cs @@ -783,21 +783,17 @@ namespace EonaCat.LogStack { try { - // Give the consumer a very short window (1 second) to finish naturally - // If it doesn't complete in time, we proceed anyway to not block app shutdown - // This is critical for ensuring the logger never blocks application exit - using (var cts = new CancellationTokenSource(TimeSpan.FromSeconds(1))) + // Bound shutdown time even if a flow is blocked while consuming. + var completed = await Task.WhenAny( + _asyncConsumer, + Task.Delay(TimeSpan.FromSeconds(1))).ConfigureAwait(false); + if (completed == _asyncConsumer) { - try - { - await _asyncConsumer.ConfigureAwait(false); - } - catch (OperationCanceledException) - { - // Timeout expired - consumer didn't finish in time - // This is acceptable; the background task will continue draining - // but we don't wait for it to allow app exit - } + await _asyncConsumer.ConfigureAwait(false); + } + else + { + _asyncCts?.Cancel(); } } catch (Exception ex) @@ -830,9 +826,13 @@ namespace EonaCat.LogStack { try { - using (var cts = new CancellationTokenSource(TimeSpan.FromSeconds(1))) + var allDisposals = Task.WhenAll(disposeTasks); + var completed = await Task.WhenAny( + allDisposals, + Task.Delay(TimeSpan.FromSeconds(1))).ConfigureAwait(false); + if (completed == allDisposals) { - await Task.WhenAll(disposeTasks).ConfigureAwait(false); + await allDisposals.ConfigureAwait(false); } } catch (OperationCanceledException) diff --git a/EonaCat.LogStack/EonaCatLoggerCore/Batching/BatchProcessor.cs b/EonaCat.LogStack/EonaCatLoggerCore/Batching/BatchProcessor.cs index 2deb31d..52f14fc 100644 --- a/EonaCat.LogStack/EonaCatLoggerCore/Batching/BatchProcessor.cs +++ b/EonaCat.LogStack/EonaCatLoggerCore/Batching/BatchProcessor.cs @@ -204,8 +204,6 @@ public sealed class BackpressureHandler : IDisposable break; } - // Collect garbage and wait - GC.Collect(GC.MaxGeneration, GCCollectionMode.Optimized, false); await Task.Delay(100, cancellationToken).ConfigureAwait(false); } } diff --git a/EonaCat.LogStack/Server.cs b/EonaCat.LogStack/Server.cs index 5041c9d..eeb18d8 100644 --- a/EonaCat.LogStack/Server.cs +++ b/EonaCat.LogStack/Server.cs @@ -79,6 +79,12 @@ namespace EonaCat.LogStack.Server /// Maximum total size of the logs directory before oldest dirs are purged (default: 10 GB). public long MaxLogDirectorySize { get; set; } = 10L * 1024 * 1024 * 1024; + /// Maximum concurrent TCP and HTTP clients (default: 128). + public int MaxConcurrentConnections { get; set; } = 128; + + /// Maximum characters accepted in one TCP or HTTP message (default: 1 MiB). + public int MaxMessageChars { get; set; } = 1024 * 1024; + /// /// Max messages accepted per second per remote endpoint (0 = disabled). /// Excess messages are dropped and counted in . @@ -98,6 +104,13 @@ namespace EonaCat.LogStack.Server private CancellationTokenSource _cts = new(); private bool _isRunning; private readonly ServerOptions _options; + private readonly SemaphoreSlim _connectionSlots; + private readonly ConcurrentDictionary _activeTcpClients = new(); + private readonly object _cleanupLock = new(); + private DateTime _nextCleanupUtc = DateTime.MinValue; + private bool _cleanupRunning; + private readonly object _rateLimiterCleanupLock = new(); + private DateTime _lastRateLimiterCleanupUtc = DateTime.MinValue; private readonly ConcurrentDictionary _rateLimiter = new(); @@ -116,7 +129,8 @@ namespace EonaCat.LogStack.Server /// Create a server using explicit . public Server(ServerOptions options) { - _options = options; + _options = options ?? throw new ArgumentNullException(nameof(options)); + _connectionSlots = new SemaphoreSlim(Math.Max(1, options.MaxConcurrentConnections)); } /// @@ -124,14 +138,14 @@ namespace EonaCat.LogStack.Server /// When is true UDP is enabled; TCP is also enabled by default. /// public Server(bool useUdp = true, int logRetentionDays = 30, long maxLogDirectorySize = 10L * 1024 * 1024 * 1024) - { - _options = new ServerOptions + : this(new ServerOptions { UseUdp = useUdp, UseTcp = true, LogRetentionDays = logRetentionDays, MaxLogDirectorySize = maxLogDirectorySize - }; + }) + { } /// Start all configured transports and begin accepting log entries. @@ -236,6 +250,10 @@ namespace EonaCat.LogStack.Server _tcpListener?.Stop(); _udpListener?.Close(); _httpListener?.Stop(); + foreach (var client in _activeTcpClients.Keys) + { + try { client.Close(); } catch { } + } _cts.Dispose(); _isRunning = false; @@ -252,15 +270,28 @@ namespace EonaCat.LogStack.Server { while (!ct.IsCancellationRequested) { - var client = await _tcpListener!.AcceptTcpClientAsync().ConfigureAwait(false); + await _connectionSlots.WaitAsync(ct).ConfigureAwait(false); + TcpClient client; + try + { + client = await _tcpListener!.AcceptTcpClientAsync().ConfigureAwait(false); + } + catch + { + _connectionSlots.Release(); + throw; + } + + _activeTcpClients.TryAdd(client, 0); Interlocked.Increment(ref Metrics.ActiveTcpConnections); - _ = Task.Run(() => HandleTcpClientAsync(client, ct), ct); + _ = HandleTcpClientAsync(client, ct); } } catch (OperationCanceledException) { // Do nothing } + catch (SocketException) when (ct.IsCancellationRequested) { } catch (Exception ex) { Console.WriteLine($"[TCP] Fatal: {ex.Message}"); } finally { Console.WriteLine("[EonaCat Server] TCP listener stopped."); } } @@ -274,27 +305,29 @@ namespace EonaCat.LogStack.Server using (var stream = client.GetStream()) using (var reader = new StreamReader(stream, Encoding.UTF8)) { - var buffer = new char[TcpReadBufferChars]; - var sb = new StringBuilder(); - int bytesRead; - - while ((bytesRead = await reader.ReadAsync(buffer, 0, buffer.Length).ConfigureAwait(false)) > 0) + var raw = await ReadBoundedMessageAsync(reader, _options.MaxMessageChars).ConfigureAwait(false); + if (raw == null) { - sb.Append(buffer, 0, bytesRead); + Interlocked.Increment(ref Metrics.TotalDropped); + return; } - var raw = sb.ToString(); Interlocked.Add(ref Metrics.TotalBytes, Encoding.UTF8.GetByteCount(raw)); await DispatchAsync(raw, remote).ConfigureAwait(false); } } - catch (Exception ex) when (!ct.IsCancellationRequested) + catch (Exception ex) { - Console.WriteLine($"[TCP] {remote}: {ex.Message}"); + if (!ct.IsCancellationRequested) + { + Console.WriteLine($"[TCP] {remote}: {ex.Message}"); + } } finally { + _activeTcpClients.TryRemove(client, out _); Interlocked.Decrement(ref Metrics.ActiveTcpConnections); + _connectionSlots.Release(); } } @@ -325,10 +358,22 @@ namespace EonaCat.LogStack.Server { while (!ct.IsCancellationRequested) { - var context = await _httpListener!.GetContextAsync().ConfigureAwait(false); - _ = Task.Run(() => HandleHttpContextAsync(context, ct), ct); + await _connectionSlots.WaitAsync(ct).ConfigureAwait(false); + HttpListenerContext context; + try + { + context = await _httpListener!.GetContextAsync().ConfigureAwait(false); + } + catch + { + _connectionSlots.Release(); + throw; + } + + _ = HandleHttpContextAsync(context, ct); } } + catch (OperationCanceledException) when (ct.IsCancellationRequested) { } catch (HttpListenerException) when (ct.IsCancellationRequested) { } catch (Exception ex) { Console.WriteLine($"[HTTP] Fatal: {ex.Message}"); } finally { Console.WriteLine("[EonaCat Server] HTTP listener stopped."); } @@ -364,7 +409,14 @@ namespace EonaCat.LogStack.Server if (req.HttpMethod == "POST" && path == "/ingest") { using var reader = new StreamReader(req.InputStream, req.ContentEncoding); - var body = await reader.ReadToEndAsync().ConfigureAwait(false); + var body = await ReadBoundedMessageAsync(reader, _options.MaxMessageChars).ConfigureAwait(false); + if (body == null) + { + Interlocked.Increment(ref Metrics.TotalDropped); + await WriteJsonResponseAsync(res, "{\"accepted\":false}", 413, ct).ConfigureAwait(false); + return; + } + Interlocked.Add(ref Metrics.TotalBytes, Encoding.UTF8.GetByteCount(body)); await DispatchAsync(body, remote).ConfigureAwait(false); await WriteJsonResponseAsync(res, "{\"accepted\":true}", 202, ct).ConfigureAwait(false); @@ -388,9 +440,30 @@ namespace EonaCat.LogStack.Server finally { try { res.Close(); } catch { } + _connectionSlots.Release(); } } + private static async Task ReadBoundedMessageAsync(TextReader reader, int maxChars) + { + maxChars = Math.Max(1, maxChars); + var buffer = new char[Math.Min(TcpReadBufferChars, maxChars)]; + var message = new StringBuilder(Math.Min(TcpReadBufferChars, maxChars)); + int charsRead; + + while ((charsRead = await reader.ReadAsync(buffer, 0, buffer.Length).ConfigureAwait(false)) > 0) + { + if (message.Length > maxChars - charsRead) + { + return null; + } + + message.Append(buffer, 0, charsRead); + } + + return message.ToString(); + } + private static async Task WriteJsonResponseAsync(HttpListenerResponse res, string json, int statusCode, CancellationToken ct) { var bytes = Encoding.UTF8.GetBytes(json); @@ -536,6 +609,7 @@ namespace EonaCat.LogStack.Server private bool IsRateLimited(string remote) { var now = DateTime.UtcNow; + PruneRateLimiter(now); var entry = _rateLimiter.GetOrAdd(remote, _ => (now, 0)); if ((now - entry.window).TotalSeconds >= 1.0) @@ -553,6 +627,29 @@ namespace EonaCat.LogStack.Server return false; } + private void PruneRateLimiter(DateTime now) + { + lock (_rateLimiterCleanupLock) + { + if (now - _lastRateLimiterCleanupUtc < TimeSpan.FromMinutes(1)) + { + return; + } + + _lastRateLimiterCleanupUtc = now; + } + + var cutoff = now.AddMinutes(-1); + var entries = (ICollection>)_rateLimiter; + foreach (var entry in _rateLimiter) + { + if (entry.Value.window < cutoff) + { + entries.Remove(entry); + } + } + } + protected virtual async Task ProcessLogAsync(string logData) { const int maxRetries = 3; @@ -605,8 +702,7 @@ namespace EonaCat.LogStack.Server Interlocked.Increment(ref Metrics.TotalWritten); LogWritten?.Invoke(logData); - // Cleanup happens asynchronously to avoid blocking - _ = Task.Run(() => CleanUpOldLogs()); + ScheduleLogCleanup(); return; } catch (UnauthorizedAccessException ex) when (attempt < maxRetries - 1) @@ -652,6 +748,40 @@ namespace EonaCat.LogStack.Server } } + private void ScheduleLogCleanup() + { + lock (_cleanupLock) + { + var now = DateTime.UtcNow; + if (_cleanupRunning || now < _nextCleanupUtc) + { + return; + } + + _cleanupRunning = true; + _nextCleanupUtc = now.AddMinutes(5); + } + + _ = Task.Run(() => + { + try + { + CleanUpOldLogs(); + } + catch (Exception ex) + { + Console.WriteLine($"[Retention] Cleanup failed: {ex.Message}"); + } + finally + { + lock (_cleanupLock) + { + _cleanupRunning = false; + } + } + }); + } + private void CleanUpOldLogs() { var root = _options.LogsRootDirectory; diff --git a/EonaCat.LogStack/Telemetry/Histogram.cs b/EonaCat.LogStack/Telemetry/Histogram.cs index dee8651..e167140 100644 --- a/EonaCat.LogStack/Telemetry/Histogram.cs +++ b/EonaCat.LogStack/Telemetry/Histogram.cs @@ -1,13 +1,22 @@ using System.Collections.Concurrent; using System.Linq; +using System.Threading; namespace EonaCat.LogStack.Telemetry; public sealed class Histogram { private readonly ConcurrentQueue _values = new(); - public void Record(double value){ _values.Enqueue(value); while(_values.Count>10000) + private int _count; + + public void Record(double value) + { + _values.Enqueue(value); + Interlocked.Increment(ref _count); + + while (Volatile.Read(ref _count) > 10000 && _values.TryDequeue(out _)) { - _values.TryDequeue(out _); + Interlocked.Decrement(ref _count); } } + public double Percentile(double p){ var a=_values.ToArray(); if(a.Length==0) { return 0; } System.Array.Sort(a); return a[(int)((a.Length-1)*p)]; } } diff --git a/EonaCat.LogStack/Telemetry/TelemetryServer.cs b/EonaCat.LogStack/Telemetry/TelemetryServer.cs index 050c0a4..36d69a9 100644 --- a/EonaCat.LogStack/Telemetry/TelemetryServer.cs +++ b/EonaCat.LogStack/Telemetry/TelemetryServer.cs @@ -9,7 +9,10 @@ namespace EonaCat.LogStack.Telemetry; public sealed class TelemetryServer : IDisposable { private const int MaxEvents = 5000; + private const int MaxConcurrentRequests = 32; private readonly HttpListener _listener = new(); + private readonly SemaphoreSlim _requestSlots = new(MaxConcurrentRequests); + private int _disposed; public ConcurrentQueue Events { get; } = new(); public int Count => Events.Count; public TelemetryServer(string prefix = "http://localhost:5155/") @@ -19,34 +22,88 @@ public sealed class TelemetryServer : IDisposable public async Task StartAsync(CancellationToken token = default) { _listener.Start(); + using var registration = token.Register(() => _listener.Stop()); while (!token.IsCancellationRequested) { - var ctx = await _listener.GetContextAsync(); - _ = Task.Run(async () => + await _requestSlots.WaitAsync(token).ConfigureAwait(false); + HttpListenerContext context; + try { - if (ctx.Request.HttpMethod == "POST" && ctx.Request.Url?.AbsolutePath == "/telemetry") - { - using var r = new StreamReader(ctx.Request.InputStream); - var item = JsonSerializer.Deserialize(await r.ReadToEndAsync()); - if (item != null) - { - Events.Enqueue(item); - // Drop the oldest events once capacity is exceeded to avoid unbounded growth. - while (Events.Count > MaxEvents && Events.TryDequeue(out _)) - { - } - } - } - else if (ctx.Request.Url?.AbsolutePath == "/") - { - var html = TelemetryDashboard.Html; - var bytes = System.Text.Encoding.UTF8.GetBytes(html); - ctx.Response.ContentType = "text/html"; - await ctx.Response.OutputStream.WriteAsync(bytes, 0, bytes.Length); - } - ctx.Response.Close(); - }); + context = await _listener.GetContextAsync().ConfigureAwait(false); + } + catch (HttpListenerException) when (token.IsCancellationRequested || Volatile.Read(ref _disposed) != 0) + { + _requestSlots.Release(); + break; + } + catch (ObjectDisposedException) when (token.IsCancellationRequested || Volatile.Read(ref _disposed) != 0) + { + _requestSlots.Release(); + break; + } + catch + { + _requestSlots.Release(); + throw; + } + + _ = ProcessRequestAsync(context); + } + } + + private async Task ProcessRequestAsync(HttpListenerContext context) + { + try + { + if (context.Request.HttpMethod == "POST" && context.Request.Url?.AbsolutePath == "/telemetry") + { + using var reader = new StreamReader(context.Request.InputStream); + var item = JsonSerializer.Deserialize(await reader.ReadToEndAsync().ConfigureAwait(false)); + if (item != null) + { + Events.Enqueue(item); + while (Events.Count > MaxEvents && Events.TryDequeue(out _)) + { + } + } + } + else if (context.Request.Url?.AbsolutePath == "/") + { + var bytes = System.Text.Encoding.UTF8.GetBytes(TelemetryDashboard.Html); + context.Response.ContentType = "text/html"; + await context.Response.OutputStream.WriteAsync(bytes, 0, bytes.Length).ConfigureAwait(false); + } + } + catch (Exception) when (Volatile.Read(ref _disposed) != 0) + { + } + catch (HttpListenerException) + { + } + catch (IOException) + { + } + catch (JsonException) + { + } + finally + { + try + { + context.Response.Close(); + } + finally + { + _requestSlots.Release(); + } + } + } + + public void Dispose() + { + if (Interlocked.Exchange(ref _disposed, 1) == 0) + { + _listener.Close(); } } - public void Dispose() => _listener.Close(); }