using Dapper; using GB5Shared.Connection; using Microsoft.Data.SqlClient; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using System.Threading.Channels; namespace GB5Shared.Telemetry.Observability; /// /// Singleton that drains a bounded /// of entries and /// bulk-inserts them into TOBSERVABILITY in batches of up to 50 rows. /// /// Write path: RequestTracingMiddleware calls (synchronous, /// never blocks). The background loop wakes on each write, accumulates up to 50 records /// or until the channel is temporarily empty, then flushes with a single Dapper call. /// /// /// Uses to obtain the /// system-level connection string — no LoginDTO required. /// /// /// Register via AddGB5Telemetry() which calls /// services.AddSingleton<IObservabilityWriter, ObservabilityWriter>() /// and services.AddHostedService(sp => sp.GetRequiredService<ObservabilityWriter>()). /// /// public sealed class ObservabilityWriter : IObservabilityWriter, IHostedService, IDisposable { // Bounded capacity of 2000: at ~1 req/ms that is 2 seconds of back-pressure // before entries are dropped. DropOldest keeps the most-recent observations. private readonly Channel _channel = Channel.CreateBounded( new BoundedChannelOptions(2000) { FullMode = BoundedChannelFullMode.DropOldest, SingleReader = true, SingleWriter = false, AllowSynchronousContinuations = false, }); private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private readonly CancellationTokenSource _cts = new(); private Task _drainTask = Task.CompletedTask; // SLA threshold: requests taking longer than this are flagged IsLowPerformance. // Configurable via appsettings: ObservabilityWriter:SlowRequestThresholdMs public int SlowRequestThresholdMs { get; set; } = 3000; public ObservabilityWriter( IServiceScopeFactory scopeFactory, ILogger logger) { _scopeFactory = scopeFactory; _logger = logger; } // ── IObservabilityWriter ───────────────────────────────────────────────── public void Enqueue(ObservabilityRecord record) => _channel.Writer.TryWrite(record); // ── IHostedService ─────────────────────────────────────────────────────── public Task StartAsync(CancellationToken cancellationToken) { _drainTask = DrainAsync(_cts.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { _channel.Writer.TryComplete(); _cts.Cancel(); try { await _drainTask.ConfigureAwait(false); } catch (OperationCanceledException) { } } // ── Background drain loop ──────────────────────────────────────────────── private async Task DrainAsync(CancellationToken ct) { var batch = new List(50); await foreach (var record in _channel.Reader.ReadAllAsync(ct).ConfigureAwait(false)) { batch.Add(record); // Drain everything currently queued (up to batch ceiling) while (batch.Count < 50 && _channel.Reader.TryRead(out var extra)) batch.Add(extra); await FlushBatchAsync(batch, ct).ConfigureAwait(false); batch.Clear(); } // Flush any remainder after channel completes if (batch.Count > 0) await FlushBatchAsync(batch, CancellationToken.None).ConfigureAwait(false); } private async Task FlushBatchAsync(List batch, CancellationToken ct) { try { await using var scope = _scopeFactory.CreateAsyncScope(); var appConnection = scope.ServiceProvider.GetRequiredService(); var connStr = await appConnection.Gb5SystemConnectionString().ConfigureAwait(false); await using var conn = new SqlConnection(connStr); await conn.OpenAsync(ct).ConfigureAwait(false); // Single parameterised INSERT per batch — Dapper expands the list const string sql = @" INSERT INTO TOBSERVABILITY ( TRACECONTEXT, CORRELATIONID, SERVICENAME, OPERATIONNAME, EVENTTYPEID, DURATIONMS, STATUSCODE, ISLOWPERFORMANCE, ISERROR, ERRORTYPE, REQUESTSIZE, RESPONSESIZE, TENANTID, USERID, SESSIONID, DIAGNOSTICSESSIONID, CREATEDON ) VALUES ( @TraceContext, @CorrelationId, @ServiceName, @OperationName, @EventTypeId, @DurationMs, @StatusCode, @IsLowPerformance, @IsError, @ErrorType, @RequestSize, @ResponseSize, @TenantId, @UserId, @SessionId, @DiagnosticSessionId, @CreatedOn )"; await conn.ExecuteAsync( new CommandDefinition(sql, batch, cancellationToken: ct)) .ConfigureAwait(false); } catch (OperationCanceledException) { } catch (Exception ex) { // Observability writes must never affect request processing _logger.LogWarning(ex, "ObservabilityWriter: flush of {Count} records failed", batch.Count); } } public void Dispose() { _cts.Dispose(); } }