using System.Collections.Concurrent; using System.IO; using System.Net; using System.Net.Sockets; using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using System; using System.Linq; using System.Collections.Generic; namespace EonaCat.LogStack.Server { // This file is part of the EonaCat project(s) which is released under the Apache License. // See the LICENSE file or go to https://EonaCat.com/License for full license details. /// /// Minimum log level the server will persist. Entries below this level are silently dropped. /// public enum ServerLogLevel { Trace = 0, Debug = 1, Information = 2, Warning = 3, Error = 4, Critical = 5 } /// /// Live counters exposed via . /// public class ServerMetrics { public long TotalReceived; public long TotalWritten; public long TotalDropped; public long TotalBytes; public long ActiveTcpConnections; public DateTime StartedAt = DateTime.UtcNow; public TimeSpan Uptime => DateTime.UtcNow - StartedAt; } /// /// Full configuration for the log server. All fields have sensible defaults. /// public class ServerOptions { /// Enable UDP listener (default: true). public bool UseUdp { get; set; } = true; /// /// Enable TCP listener (default: true). /// Both TCP and UDP can run simultaneously on the same port. /// public bool UseTcp { get; set; } = true; /// /// Enable a minimal HTTP ingest endpoint (default: false). /// POST /ingest - accepts a single JSON entry or a JSON array. /// GET /metrics - returns a JSON metrics snapshot. /// public bool UseHttp { get; set; } = false; /// Port for UDP/TCP listeners (default: 5555). public int Port { get; set; } = 5555; /// Port for the HTTP ingest endpoint (default: 5556). public int HttpPort { get; set; } = 5556; /// IP address to bind to (null = IPAddress.Any). public IPAddress? BindAddress { get; set; } /// /// Only persist entries at or above this level (default: Trace = everything). /// Plain-text entries (non-JSON) always pass the level filter. /// public ServerLogLevel MinimumLevel { get; set; } = ServerLogLevel.Trace; /// Max log file size before rolling over (default: 200 MB). public long MaxLogFileSize { get; set; } = 200L * 1024 * 1024; /// Days to keep daily log directories (default: 30). public int LogRetentionDays { get; set; } = 30; /// 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 . /// public int RateLimitPerSecond { get; set; } = 0; /// Root directory where log files are written (default: "logs"). public string LogsRootDirectory { get; set; } = "logs"; } public class Server { private TcpListener? _tcpListener; private UdpClient? _udpListener; private HttpListener? _httpListener; 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(); private const int UdpBufferSize = 65507; private const int TcpReadBufferChars = 8192; /// Live connection and throughput counters. public ServerMetrics Metrics { get; } = new(); /// Raised after a log line has been written to disk. public event Action? LogWritten; /// Raised when a log line is dropped (level filter or rate limit). public event Action? LogDropped; /// Create a server using explicit . public Server(ServerOptions options) { _options = options ?? throw new ArgumentNullException(nameof(options)); _connectionSlots = new SemaphoreSlim(Math.Max(1, options.MaxConcurrentConnections)); } /// /// Backwards-compatible constructor matching the original API. /// 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) : this(new ServerOptions { UseUdp = useUdp, UseTcp = true, LogRetentionDays = logRetentionDays, MaxLogDirectorySize = maxLogDirectorySize }) { } /// Start all configured transports and begin accepting log entries. public async Task Start(IPAddress? ipAddress = null, int port = 5555) { if (_isRunning) { Console.WriteLine("[EonaCat Server] Already running."); return; } _cts = new CancellationTokenSource(); _isRunning = true; Metrics.StartedAt = DateTime.UtcNow; // Initialize log directory with permission checking if (!InitializeLogDirectory()) { Console.WriteLine("[EonaCat Server] WARNING: Could not initialize log directory. Logs will be dropped if write fails."); } var bind = ipAddress ?? _options.BindAddress ?? IPAddress.Any; var tasks = new List(); if (_options.UseTcp) { _tcpListener = new TcpListener(bind, port); _tcpListener.Start(); Console.WriteLine($"[EonaCat Server] TCP listener on {bind}:{port}"); tasks.Add(ListenTcpAsync(_cts.Token)); } if (_options.UseUdp) { _udpListener = new UdpClient(port); Console.WriteLine($"[EonaCat Server] UDP listener on port {port}"); tasks.Add(ListenUdpAsync(_cts.Token)); } if (_options.UseHttp) { _httpListener = new HttpListener(); _httpListener.Prefixes.Add($"http://*:{_options.HttpPort}/"); _httpListener.Start(); Console.WriteLine($"[EonaCat Server] HTTP ingest on port {_options.HttpPort}"); tasks.Add(ListenHttpAsync(_cts.Token)); } if (tasks.Count == 0) { Console.WriteLine("[EonaCat Server] WARNING: No transports are enabled. Check ServerOptions."); _isRunning = false; return; } await Task.WhenAll(tasks).ConfigureAwait(false); } /// /// Initializes the log directory during startup, ensuring write permissions. /// Returns true if successful, false if there were issues (but doesn't prevent startup). /// private bool InitializeLogDirectory() { try { var logDir = _options.LogsRootDirectory; if (!Helpers.DirectoryPermissionHelper.EnsureDirectoryHierarchy(logDir)) { Console.WriteLine($"[EonaCat Server] WARNING: Could not fully initialize log directory: {logDir}"); Console.WriteLine($"[EonaCat Server] Diagnosis: {Helpers.DirectoryPermissionHelper.GetAccessIssueDiagnosis(logDir)}"); return false; } if (!Helpers.DirectoryPermissionHelper.CanWrite(logDir)) { Console.WriteLine($"[EonaCat Server] WARNING: Log directory exists but is not writable: {logDir}"); Console.WriteLine($"[EonaCat Server] Please check file system permissions."); return false; } Console.WriteLine($"[EonaCat Server] Log directory ready: {Path.GetFullPath(logDir)}"); return true; } catch (Exception ex) { Console.WriteLine($"[EonaCat Server] Error initializing log directory: {ex.GetType().Name}: {ex.Message}"); return false; } } /// Gracefully stop all transports. public void Stop() { if (!_isRunning) { return; } _cts.Cancel(); _tcpListener?.Stop(); _udpListener?.Close(); _httpListener?.Stop(); foreach (var client in _activeTcpClients.Keys) { try { client.Close(); } catch { } } _cts.Dispose(); _isRunning = false; Console.WriteLine( $"[EonaCat Server] Stopped. " + $"Received={Metrics.TotalReceived} Written={Metrics.TotalWritten} " + $"Dropped={Metrics.TotalDropped} Bytes={Metrics.TotalBytes} " + $"Uptime={Metrics.Uptime:hh\\:mm\\:ss}"); } private async Task ListenTcpAsync(CancellationToken ct) { try { while (!ct.IsCancellationRequested) { 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); _ = 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."); } } private async Task HandleTcpClientAsync(TcpClient client, CancellationToken ct) { var remote = client.Client.RemoteEndPoint?.ToString() ?? "unknown"; try { using (client) using (var stream = client.GetStream()) using (var reader = new StreamReader(stream, Encoding.UTF8)) { var raw = await ReadBoundedMessageAsync(reader, _options.MaxMessageChars).ConfigureAwait(false); if (raw == null) { Interlocked.Increment(ref Metrics.TotalDropped); return; } Interlocked.Add(ref Metrics.TotalBytes, Encoding.UTF8.GetByteCount(raw)); await DispatchAsync(raw, remote).ConfigureAwait(false); } } catch (Exception ex) { if (!ct.IsCancellationRequested) { Console.WriteLine($"[TCP] {remote}: {ex.Message}"); } } finally { _activeTcpClients.TryRemove(client, out _); Interlocked.Decrement(ref Metrics.ActiveTcpConnections); _connectionSlots.Release(); } } private async Task ListenUdpAsync(CancellationToken ct) { try { while (!ct.IsCancellationRequested) { var result = await _udpListener!.ReceiveAsync().ConfigureAwait(false); var remote = result.RemoteEndPoint.ToString(); Interlocked.Add(ref Metrics.TotalBytes, result.Buffer.Length); var raw = Encoding.UTF8.GetString(result.Buffer); await DispatchAsync(raw, remote).ConfigureAwait(false); } } catch (OperationCanceledException) { // Do nothing } catch (Exception ex) { Console.WriteLine($"[UDP] Fatal: {ex.Message}"); } finally { Console.WriteLine("[EonaCat Server] UDP listener stopped."); } } private async Task ListenHttpAsync(CancellationToken ct) { try { while (!ct.IsCancellationRequested) { 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."); } } private async Task HandleHttpContextAsync(HttpListenerContext ctx, CancellationToken ct) { var req = ctx.Request; var res = ctx.Response; var remote = req.RemoteEndPoint?.ToString() ?? "http"; var path = req.Url?.AbsolutePath.TrimEnd('/') ?? ""; try { // GET /metrics if (req.HttpMethod == "GET" && path == "/metrics") { var json = JsonSerializer.Serialize(new { totalReceived = Metrics.TotalReceived, totalWritten = Metrics.TotalWritten, totalDropped = Metrics.TotalDropped, totalBytes = Metrics.TotalBytes, activeTcpConnections = Metrics.ActiveTcpConnections, uptimeSeconds = (long)Metrics.Uptime.TotalSeconds, startedAt = Metrics.StartedAt.ToString("o") }); await WriteJsonResponseAsync(res, json, 200, ct).ConfigureAwait(false); return; } // POST /ingest if (req.HttpMethod == "POST" && path == "/ingest") { using var reader = new StreamReader(req.InputStream, req.ContentEncoding); 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); return; } res.StatusCode = 404; } catch (Exception ex) { Console.WriteLine($"[HTTP] Handler error: {ex.Message}"); try { res.StatusCode = 500; } catch { // Do nothing } } 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); res.StatusCode = statusCode; res.ContentType = "application/json; charset=utf-8"; res.ContentLength64 = bytes.Length; await res.OutputStream.WriteAsync(bytes, 0, bytes.Length, ct).ConfigureAwait(false); } private async Task DispatchAsync(string raw, string remote) { if (string.IsNullOrWhiteSpace(raw)) { return; } Interlocked.Increment(ref Metrics.TotalReceived); // Rate limiting if (_options.RateLimitPerSecond > 0 && IsRateLimited(remote)) { Interlocked.Increment(ref Metrics.TotalDropped); LogDropped?.Invoke(raw); return; } // Try structured JSON path first var entries = TryParseStructured(raw); if (entries.Count == 0) { // Plain text - always passes level filter await ProcessLogAsync(raw).ConfigureAwait(false); return; } foreach (var (line, level) in entries) { if (level < _options.MinimumLevel) { Interlocked.Increment(ref Metrics.TotalDropped); LogDropped?.Invoke(line); continue; } await ProcessLogAsync(line).ConfigureAwait(false); } } private List<(string line, ServerLogLevel level)> TryParseStructured(string raw) { var trimmed = raw.Trim(); if (trimmed.Length == 0 || (trimmed[0] != '{' && trimmed[0] != '[')) { return new List<(string, ServerLogLevel)>(); } var result = new List<(string, ServerLogLevel)>(); try { using var doc = JsonDocument.Parse(trimmed); if (doc.RootElement.ValueKind == JsonValueKind.Array) { foreach (var el in doc.RootElement.EnumerateArray()) { result.Add(ParseJsonEntry(el)); } } else { result.Add(ParseJsonEntry(doc.RootElement)); } } catch { // if parsing fails, treat as plain text } return result; } private static (string line, ServerLogLevel level) ParseJsonEntry(JsonElement jsonElement) { var level = ExtractLevel(jsonElement); var timestamp = jsonElement.TryGetProperty("timestamp", out var tsProp) ? tsProp.GetString() ?? "" : DateTime.UtcNow.ToString("o"); var message = jsonElement.TryGetProperty("message", out var mProp) ? mProp.GetString() ?? "" : jsonElement.TryGetProperty("Message", out var mProp2) ? mProp2.GetString() ?? "" : jsonElement.GetRawText(); var source = jsonElement.TryGetProperty("source", out var sProp) ? sProp.GetString() ?? "" : jsonElement.TryGetProperty("application", out var aProp) ? aProp.GetString() ?? "" : ""; var host = jsonElement.TryGetProperty("host", out var hProp) ? hProp.GetString() ?? "" : ""; var exception = jsonElement.TryGetProperty("exception", out var eProp) ? eProp.GetString() : null; var traceId = jsonElement.TryGetProperty("traceId", out var trProp) ? trProp.GetString() : null; var stringBuilder = new StringBuilder(); stringBuilder.Append($"[{timestamp}] [{level.ToString().ToUpper()}]"); if (!string.IsNullOrEmpty(source)) { stringBuilder.Append($" [{source}]"); } if (!string.IsNullOrEmpty(host)) { stringBuilder.Append($" host={host}"); } if (!string.IsNullOrEmpty(traceId)) { stringBuilder.Append($" trace={traceId}"); } stringBuilder.Append($" {message}"); if (!string.IsNullOrEmpty(exception)) { stringBuilder.Append($"\n EXCEPTION: {exception}"); } return (stringBuilder.ToString(), level); } private static ServerLogLevel ExtractLevel(JsonElement element) { string? level = null; foreach (var key in new[] { "level", "Level", "severity", "Severity" }) { if (element.TryGetProperty(key, out var p)) { level = p.GetString(); break; } } return level?.ToLowerInvariant() switch { "trace" or "verbose" => ServerLogLevel.Trace, "debug" => ServerLogLevel.Debug, "info" or "information" => ServerLogLevel.Information, "warn" or "warning" => ServerLogLevel.Warning, "error" => ServerLogLevel.Error, "fatal" or "critical" => ServerLogLevel.Critical, _ => ServerLogLevel.Information }; } 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) { _rateLimiter[remote] = (now, 1); return false; } if (entry.count >= _options.RateLimitPerSecond) { return true; } _rateLimiter[remote] = (entry.window, entry.count + 1); 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; const int retryDelayMs = 100; for (int attempt = 0; attempt < maxRetries; attempt++) { try { var root = _options.LogsRootDirectory; // Ensure directory exists with permission handling if (!EnsureLogDirectory(root)) { if (attempt < maxRetries - 1) { await Task.Delay(retryDelayMs).ConfigureAwait(false); continue; } // Last attempt failed - log that we couldn't write Console.WriteLine($"[EonaCat Server] WARNING: Cannot write to log directory. Reason: {Helpers.DirectoryPermissionHelper.GetAccessIssueDiagnosis(root)}"); Interlocked.Increment(ref Metrics.TotalDropped); return; } var daily = Path.Combine(root, DateTime.Now.ToString("yyyyMMdd")); if (!EnsureLogDirectory(daily)) { if (attempt < maxRetries - 1) { await Task.Delay(retryDelayMs).ConfigureAwait(false); continue; } Console.WriteLine($"[EonaCat Server] WARNING: Cannot create daily log directory: {daily}"); Interlocked.Increment(ref Metrics.TotalDropped); return; } var basePath = Path.Combine(daily, "EonaCatLogs"); var filePath = basePath + ".log"; int idx = 1; while (File.Exists(filePath) && new FileInfo(filePath).Length > _options.MaxLogFileSize) { filePath = $"{basePath}_{idx++}.log"; } using (var stream = new FileStream(filePath, FileMode.Append, FileAccess.Write, FileShare.Read | FileShare.Delete)) using (var writer = new StreamWriter(stream)) { writer.WriteLine(logData); } Interlocked.Increment(ref Metrics.TotalWritten); LogWritten?.Invoke(logData); ScheduleLogCleanup(); return; } catch (UnauthorizedAccessException ex) when (attempt < maxRetries - 1) { Console.WriteLine($"[EonaCat Server] Access denied to log directory (attempt {attempt + 1}/{maxRetries}): {ex.Message}"); await Task.Delay(retryDelayMs * (attempt + 1)).ConfigureAwait(false); } catch (IOException ex) when (attempt < maxRetries - 1) { Console.WriteLine($"[EonaCat Server] IO error writing to log (attempt {attempt + 1}/{maxRetries}): {ex.Message}"); await Task.Delay(retryDelayMs * (attempt + 1)).ConfigureAwait(false); } catch (Exception ex) { // Any other exception - log and drop the message to prevent crash Console.WriteLine($"[EonaCat Server] Error writing log: {ex.GetType().Name}: {ex.Message}"); Interlocked.Increment(ref Metrics.TotalDropped); LogDropped?.Invoke(logData); return; } } // All retries exhausted Console.WriteLine("[EonaCat Server] Max retries exhausted, dropping log message"); Interlocked.Increment(ref Metrics.TotalDropped); LogDropped?.Invoke(logData); } /// /// Ensures a log directory exists and is writable, with permission fixing if needed. /// Returns true if successful, false otherwise. /// private bool EnsureLogDirectory(string dirPath) { try { return Helpers.DirectoryPermissionHelper.EnsureDirectoryHierarchy(dirPath); } catch (Exception ex) { Console.WriteLine($"[EonaCat Server] Failed to ensure log directory: {ex.Message}"); return false; } } 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; if (!Directory.Exists(root)) { return; } foreach (var dir in Directory.GetDirectories(root)) { try { if (new DirectoryInfo(dir).CreationTime < DateTime.Now.AddDays(-_options.LogRetentionDays)) { Directory.Delete(dir, true); Console.WriteLine($"[Retention] Deleted expired directory: {dir}"); } } catch (Exception ex) { Console.WriteLine($"[Retention] Error: {ex.Message}"); } } long total = GetDirectorySize(root); if (total <= _options.MaxLogDirectorySize) { return; } Console.WriteLine($"[Retention] Directory {total / (1024 * 1024)} MB over limit. Purging oldest..."); foreach (var dir in Directory.GetDirectories(root).OrderBy(d => new DirectoryInfo(d).CreationTime)) { try { total -= GetDirectorySize(dir); Directory.Delete(dir, true); Console.WriteLine($"[Retention] Purged: {dir}"); if (total <= _options.MaxLogDirectorySize) { break; } } catch (Exception ex) { Console.WriteLine($"[Retention] Error purging {dir}: {ex.Message}"); } } } private static long GetDirectorySize(string directory) { long size = 0; try { size += Directory.GetFiles(directory).Sum(f => new FileInfo(f).Length); foreach (var sub in Directory.GetDirectories(directory)) { size += GetDirectorySize(sub); } } catch { // ignored } return size; } } }