This commit is contained in:
1 parent
cfe8eb74ec
commit
4081b5949d
5 files changed
+245
-51
No files matched your search
@@ -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)))
|
||||
{
|
||||
try
|
||||
// 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)
|
||||
{
|
||||
await _asyncConsumer.ConfigureAwait(false);
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
else
|
||||
{
|
||||
// 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
|
||||
}
|
||||
_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)
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
+149
-19
@@ -79,6 +79,12 @@ namespace EonaCat.LogStack.Server
|
||||
/// <summary>Maximum total size of the logs directory before oldest dirs are purged (default: 10 GB).</summary>
|
||||
public long MaxLogDirectorySize { get; set; } = 10L * 1024 * 1024 * 1024;
|
||||
|
||||
/// <summary>Maximum concurrent TCP and HTTP clients (default: 128).</summary>
|
||||
public int MaxConcurrentConnections { get; set; } = 128;
|
||||
|
||||
/// <summary>Maximum characters accepted in one TCP or HTTP message (default: 1 MiB).</summary>
|
||||
public int MaxMessageChars { get; set; } = 1024 * 1024;
|
||||
|
||||
/// <summary>
|
||||
/// Max messages accepted per second per remote endpoint (0 = disabled).
|
||||
/// Excess messages are dropped and counted in <see cref="ServerMetrics.TotalDropped"/>.
|
||||
@@ -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<TcpClient, byte> _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<string, (DateTime window, int count)> _rateLimiter = new();
|
||||
|
||||
@@ -116,7 +129,8 @@ namespace EonaCat.LogStack.Server
|
||||
/// <summary>Create a server using explicit <see cref="ServerOptions"/>.</summary>
|
||||
public Server(ServerOptions options)
|
||||
{
|
||||
_options = options;
|
||||
_options = options ?? throw new ArgumentNullException(nameof(options));
|
||||
_connectionSlots = new SemaphoreSlim(Math.Max(1, options.MaxConcurrentConnections));
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -124,14 +138,14 @@ namespace EonaCat.LogStack.Server
|
||||
/// When <paramref name="useUdp"/> is true UDP is enabled; TCP is also enabled by default.
|
||||
/// </summary>
|
||||
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
|
||||
};
|
||||
})
|
||||
{
|
||||
}
|
||||
|
||||
/// <summary>Start all configured transports and begin accepting log entries.</summary>
|
||||
@@ -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)
|
||||
{
|
||||
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<string?> 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<KeyValuePair<string, (DateTime window, int count)>>)_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;
|
||||
|
||||
@@ -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<double> _values = new();
|
||||
public void Record(double value){ _values.Enqueue(value); while(_values.Count>10000)
|
||||
private int _count;
|
||||
|
||||
public void Record(double value)
|
||||
{
|
||||
_values.TryDequeue(out _);
|
||||
_values.Enqueue(value);
|
||||
Interlocked.Increment(ref _count);
|
||||
|
||||
while (Volatile.Read(ref _count) > 10000 && _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)]; }
|
||||
}
|
||||
@@ -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<TelemetryEvent> 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")
|
||||
context = await _listener.GetContextAsync().ConfigureAwait(false);
|
||||
}
|
||||
catch (HttpListenerException) when (token.IsCancellationRequested || Volatile.Read(ref _disposed) != 0)
|
||||
{
|
||||
using var r = new StreamReader(ctx.Request.InputStream);
|
||||
var item = JsonSerializer.Deserialize<TelemetryEvent>(await r.ReadToEndAsync());
|
||||
_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<TelemetryEvent>(await reader.ReadToEndAsync().ConfigureAwait(false));
|
||||
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 == "/")
|
||||
else if (context.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();
|
||||
});
|
||||
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();
|
||||
}
|
||||
Reference in new issue
Block a user