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();
}
}