using GB5Shared.Telemetry.Observability; using Microsoft.AspNetCore.Http; using System.Diagnostics; using System.Text; namespace GB5Shared.Telemetry.Middleware; /// /// ASP.NET Core middleware that enriches the active trace span with comprehensive request, /// response, and exception context — automatically, for every API entry point (FastEndpoints, /// Controllers, Minimal API, ReportEndpoints, etc.) without touching individual files. /// /// What it captures on every request: /// /// HTTP headers: correlation ID, real client IP, Dapr source app-id, WIP approval marker. /// Request payload: JSON body as http.request.body tag (up to ). /// Response: HTTP status code and content-type after the endpoint completes. /// Exceptions: full details — type, message, stack trace, inner exception chain, trace/span IDs — when an exception escapes to this layer. /// /// /// /// Register in Program.cs after app.UseRouting() so the root span is active: /// app.UseMiddleware<RequestTracingMiddleware>(); /// /// /// is injected via DI and registered automatically by /// AddGB5Telemetry() — no manual registration required. /// /// public sealed class RequestTracingMiddleware { private readonly RequestDelegate _next; private readonly TelemetryOptions _options; private readonly IObservabilityWriter? _observabilityWriter; private readonly string _serviceName; public RequestTracingMiddleware( RequestDelegate next, TelemetryOptions options, IObservabilityWriter? observabilityWriter = null) { _next = next; _options = options; _observabilityWriter = observabilityWriter; // Real per-service name, set by AddGB5Telemetry(configuration, serviceName) in each // service's Program.cs — never a shared hardcoded literal (see TelemetryOptions.ServiceName). _serviceName = options.ServiceName; } public async Task InvokeAsync(HttpContext ctx) { var activity = Activity.Current; var startTime = Stopwatch.GetTimestamp(); if (activity is not null) { EnrichFromHeaders(activity, ctx); if (_options.CaptureRequestBody) await CaptureRequestBodyAsync(activity, ctx.Request); } PropagateCorrelationId(ctx, activity); Exception? capturedException = null; try { await _next(ctx); EnrichResponse(activity, ctx); } catch (Exception ex) { // Safety net: catches exceptions that escape FastEndpoints, Controllers, etc. // BaseEndpoint already handles its own exceptions internally, but Controllers, // Minimal API handlers, and custom endpoints may let exceptions propagate here. if (activity is not null && _options.CaptureExceptionDetails) EnrichException(activity, ex, ctx, _options); capturedException = ex; throw; } finally { // Enqueue observability record — fire-and-forget, never blocks the request if (_observabilityWriter is not null) { var durationMs = (int)Stopwatch.GetElapsedTime(startTime).TotalMilliseconds; _observabilityWriter.Enqueue(BuildRecord(ctx, activity, durationMs, capturedException)); } } } private ObservabilityRecord BuildRecord( HttpContext ctx, Activity? activity, int durationMs, Exception? ex) { var correlationId = ctx.Request.Headers["X-Correlation-Id"].FirstOrDefault() ?? activity?.TraceId.ToString(); // Read tags set by BaseEndpoint during request execution int.TryParse(activity?.GetTagItem("gb5.client.id") as string, out var clientId); int.TryParse(activity?.GetTagItem("gb5.user.id") as string, out var userId); var sessionId = activity?.GetTagItem("gb5.session.id") as string; return new ObservabilityRecord { TraceContext = activity is not null ? $"{activity.TraceId}:{activity.SpanId}" : null, CorrelationId = correlationId, ServiceName = _serviceName, OperationName = ctx.Request.Path.Value ?? "/", DurationMs = durationMs, StatusCode = ctx.Response.StatusCode, IsLowPerformance = durationMs > 3000, // 3 s default SLA IsError = ctx.Response.StatusCode >= 500 || ex is not null, ErrorType = ex?.GetType().Name, RequestSize = (int?)ctx.Request.ContentLength, ResponseSize = (int?)ctx.Response.ContentLength, TenantId = clientId == 0 ? -1 : clientId, UserId = userId == 0 ? null : userId, SessionId = sessionId, CreatedOn = DateTime.UtcNow, }; } // ── Request enrichment ──────────────────────────────────────────────────── private static void EnrichFromHeaders(Activity activity, HttpContext ctx) { var req = ctx.Request; // ── Correlation ID ──────────────────────────────────────────────────── var correlationId = req.Headers["X-Correlation-Id"].FirstOrDefault() ?? activity.TraceId.ToString(); activity.SetTag("gb5.correlation.id", correlationId); // ── Real client IP ──────────────────────────────────────────────────── var forwardedFor = req.Headers["X-Forwarded-For"].FirstOrDefault(); if (!string.IsNullOrEmpty(forwardedFor)) activity.SetTag("http.client_ip", forwardedFor.Split(',')[0].Trim()); else if (ctx.Connection.RemoteIpAddress is not null) activity.SetTag("http.client_ip", ctx.Connection.RemoteIpAddress.ToString()); // ── Dapr source identification ──────────────────────────────────────── var daprAppId = req.Headers["dapr-app-id"].FirstOrDefault(); if (!string.IsNullOrEmpty(daprAppId)) activity.SetTag("dapr.source.app_id", daprAppId); // ── WIP approval replay ─────────────────────────────────────────────── var wipHeader = req.Headers["X-Wip-Approval"].FirstOrDefault(); if (!string.IsNullOrEmpty(wipHeader)) activity.SetTag("gb5.wip.approval_id", wipHeader); // ── Request metadata ────────────────────────────────────────────────── if (req.ContentLength.HasValue) activity.SetTag("http.request_content_length", req.ContentLength.Value); activity.SetTag("http.request.method", req.Method); activity.SetTag("http.url", $"{req.Scheme}://{req.Host}{req.Path}{req.QueryString}"); } private async Task CaptureRequestBodyAsync(Activity activity, HttpRequest req) { // Only capture JSON — skip multipart/form-data, binary uploads, etc. if (req.ContentType?.Contains("application/json", StringComparison.OrdinalIgnoreCase) != true) return; // EnableBuffering converts the body to a seekable stream so the endpoint can still read it. req.EnableBuffering(); try { req.Body.Position = 0; var limit = _options.MaxRequestBodyBytes; var buffer = new byte[limit + 1]; // +1 detects truncation var total = 0; int read; while (total < buffer.Length && (read = await req.Body.ReadAsync(buffer.AsMemory(total))) > 0) total += read; req.Body.Position = 0; // Reset so the endpoint reads the full body if (total == 0) return; var truncated = total > limit; var body = Encoding.UTF8.GetString(buffer, 0, truncated ? limit : total); if (truncated) body += "...[truncated]"; activity.SetTag("http.request.body", body); } catch { // Best effort — never fail the request over a tracing operation try { req.Body.Position = 0; } catch { } } } // ── Response enrichment ─────────────────────────────────────────────────── private static void EnrichResponse(Activity? activity, HttpContext ctx) { if (activity is null) return; activity.SetTag("http.response.status_code", ctx.Response.StatusCode); if (ctx.Response.ContentType is { Length: > 0 } ct) activity.SetTag("http.response.content_type", ct); if (ctx.Response.ContentLength.HasValue) activity.SetTag("http.response.content_length", ctx.Response.ContentLength.Value); } // ── Exception enrichment ────────────────────────────────────────────────── private static void EnrichException( Activity activity, Exception ex, HttpContext ctx, TelemetryOptions options) { activity.SetStatus(ActivityStatusCode.Error, ex.Message); activity.SetTag("error", true); activity.SetTag("exception.escaped", true); activity.SetTag("exception.type", ex.GetType().FullName ?? ex.GetType().Name); activity.SetTag("exception.message", ex.Message); // Stack trace — first 15 frames to keep tag size manageable if (ex.StackTrace is { Length: > 0 } st) { var frames = st.Split('\n').Take(15); activity.SetTag("exception.stacktrace", string.Join('\n', frames)); } // Inner exception chain var inner = ex.InnerException; for (var depth = 0; inner is not null && depth < options.MaxExceptionDepth; depth++) { activity.SetTag($"exception.inner[{depth}].type", inner.GetType().Name); activity.SetTag($"exception.inner[{depth}].message", inner.Message); inner = inner.InnerException; } // Trace / correlation IDs — essential for cross-service error lookup var correlId = ctx.Request.Headers["X-Correlation-Id"].FirstOrDefault() ?? activity.TraceId.ToString(); activity.SetTag("gb5.error.correlation_id", correlId); activity.SetTag("gb5.error.trace_id", activity.TraceId.ToString()); activity.SetTag("gb5.error.span_id", activity.SpanId.ToString()); activity.SetTag("gb5.error.request_path", ctx.Request.Path.Value); activity.SetTag("gb5.error.request_method", ctx.Request.Method); // OTel standard exception event — visible in Jaeger "Logs" tab and Grafana Tempo activity.AddEvent(new ActivityEvent("exception", tags: new ActivityTagsCollection { ["exception.type"] = ex.GetType().FullName ?? ex.GetType().Name, ["exception.message"] = ex.Message, ["exception.stacktrace"] = ex.StackTrace ?? string.Empty, ["exception.escaped"] = true, })); } // ── Correlation ID propagation ──────────────────────────────────────────── private static void PropagateCorrelationId(HttpContext ctx, Activity? activity) { if (ctx.Response.Headers.ContainsKey("X-Correlation-Id")) return; var id = ctx.Request.Headers["X-Correlation-Id"].FirstOrDefault() ?? activity?.TraceId.ToString() ?? Guid.NewGuid().ToString("N"); ctx.Response.Headers["X-Correlation-Id"] = id; } }