using System.Globalization; using System.Net.Http.Headers; using System.Text.Json; using Microsoft.Extensions.Caching.Memory; namespace Jiaowu.Api.Infrastructure.Observability; public sealed class PerformanceReportService( HttpClient httpClient, IMemoryCache cache, PerformanceReportingOptions options, ObservabilityOptions observability, ILogger logger) { public async Task GetAsync( string? range, CancellationToken cancellationToken) { var rangeSpec = PerformanceRange.TryParse(range); if (rangeSpec is null) throw new ArgumentOutOfRangeException( nameof(range), "性能报表范围仅支持 15m、1h、24h 或 7d。"); if (!options.Enabled || string.IsNullOrWhiteSpace(options.PrometheusBaseUrl)) { return PerformanceReport.NotConfigured( rangeSpec.Key, options.GrafanaBaseUrl); } var cacheKey = $"performance-report:{rangeSpec.Key}"; if (cache.TryGetValue(cacheKey, out var cached)) return cached!; PerformanceReport report; try { report = await LoadAsync(rangeSpec, cancellationToken); } catch (Exception exception) when ( !cancellationToken.IsCancellationRequested) { logger.LogWarning( exception, "Performance report source is unavailable for range {Range}.", rangeSpec.Key); report = PerformanceReport.Unavailable( rangeSpec.Key, options.GrafanaBaseUrl); } cache.Set( cacheKey, report, TimeSpan.FromSeconds(options.CacheSeconds)); return report; } private async Task LoadAsync( PerformanceRange range, CancellationToken cancellationToken) { var now = DateTime.UtcNow; var from = now - range.Duration; var requestCountSelector = Selector( options.RequestDurationMetric + "_count", (options.ServiceLabel, "=", observability.ServiceName)); var requestBucketSelector = Selector( options.RequestDurationMetric + "_bucket", (options.ServiceLabel, "=", observability.ServiceName)); var errorCountSelector = Selector( options.RequestDurationMetric + "_count", (options.ServiceLabel, "=", observability.ServiceName), ("http_response_status_code", "=~", "5..")); var databaseCountSelector = Selector( options.DatabaseDurationMetric + "_count", (options.ServiceLabel, "=", observability.ServiceName)); var databaseBucketSelector = Selector( options.DatabaseDurationMetric + "_bucket", (options.ServiceLabel, "=", observability.ServiceName)); var slowDatabaseSelector = Selector( options.SlowDatabaseMetric, (options.ServiceLabel, "=", observability.ServiceName)); var failedDatabaseSelector = Selector( options.FailedDatabaseMetric, (options.ServiceLabel, "=", observability.ServiceName)); var requestCountTask = QueryScalarAsync( $"sum(increase({requestCountSelector}[{range.PrometheusRange}]))", now, cancellationToken); var serverErrorCountTask = QueryScalarAsync( $"sum(increase({errorCountSelector}[{range.PrometheusRange}]))", now, cancellationToken); var requestP95Task = QueryScalarAsync( "histogram_quantile(0.95, " + $"sum by (le) (rate({requestBucketSelector}[{range.RateWindow}]))) " + "* 1000", now, cancellationToken); var databaseP95Task = QueryScalarAsync( "histogram_quantile(0.95, " + $"sum by (le) (rate({databaseBucketSelector}[{range.RateWindow}])))", now, cancellationToken); var slowCountTask = QueryScalarAsync( $"sum(increase({slowDatabaseSelector}[{range.PrometheusRange}]))", now, cancellationToken); var failedCountTask = QueryScalarAsync( $"sum(increase({failedDatabaseSelector}[{range.PrometheusRange}]))", now, cancellationToken); var requestTimelineTask = QueryRangeAsync( $"sum(rate({requestCountSelector}[{range.RateWindow}]))", from, now, range.StepSeconds, cancellationToken); var latencyTimelineTask = QueryRangeAsync( "histogram_quantile(0.95, " + $"sum by (le) (rate({requestBucketSelector}[{range.RateWindow}]))) " + "* 1000", from, now, range.StepSeconds, cancellationToken); var routeLatencyTask = QueryVectorAsync( "histogram_quantile(0.95, " + $"sum by (le, http_route) (rate({requestBucketSelector}" + $"[{range.RateWindow}]))) * 1000", now, cancellationToken); var routeCountTask = QueryVectorAsync( $"sum by (http_route) (increase({requestCountSelector}" + $"[{range.PrometheusRange}]))", now, cancellationToken); var routeErrorTask = QueryVectorAsync( $"sum by (http_route) (increase({errorCountSelector}" + $"[{range.PrometheusRange}]))", now, cancellationToken); var queryLatencyTask = QueryVectorAsync( "histogram_quantile(0.95, " + $"sum by (le, db_query_name) (rate({databaseBucketSelector}" + $"[{range.RateWindow}])))", now, cancellationToken); var queryCountTask = QueryVectorAsync( $"sum by (db_query_name) (increase({databaseCountSelector}" + $"[{range.PrometheusRange}]))", now, cancellationToken); var querySlowTask = QueryVectorAsync( $"sum by (db_query_name) (increase({slowDatabaseSelector}" + $"[{range.PrometheusRange}]))", now, cancellationToken); await Task.WhenAll( requestCountTask, serverErrorCountTask, requestP95Task, databaseP95Task, slowCountTask, failedCountTask, requestTimelineTask, latencyTimelineTask, routeLatencyTask, routeCountTask, routeErrorTask, queryLatencyTask, queryCountTask, querySlowTask); var requestCount = await requestCountTask; var serverErrorCount = await serverErrorCountTask; double? errorRate = requestCount is > 0 && serverErrorCount.HasValue ? serverErrorCount.Value / requestCount.Value * 100 : requestCount == 0 ? 0 : null; var timeline = MergeTimeline( await requestTimelineTask, await latencyTimelineTask); var endpoints = MergeRanking( await routeLatencyTask, await routeCountTask, await routeErrorTask, "http_route"); var databaseQueries = MergeRanking( await queryLatencyTask, await queryCountTask, await querySlowTask, "db_query_name"); return new PerformanceReport( "ready", range.Key, from, now, DateTime.UtcNow, "prometheus", EmptyToNull(options.GrafanaBaseUrl), null, new PerformanceHeadline( Round(requestCount), Round(await requestP95Task), Round(errorRate), Round(await databaseP95Task), Round(await slowCountTask), Round(await failedCountTask)), timeline, endpoints, databaseQueries); } private async Task QueryScalarAsync( string query, DateTime time, CancellationToken cancellationToken) { var vector = await QueryVectorAsync( query, time, cancellationToken); return vector.FirstOrDefault()?.Value; } private async Task> QueryVectorAsync( string query, DateTime time, CancellationToken cancellationToken) { var uri = BuildUri( "api/v1/query", ("query", query), ("time", ToUnixSeconds(time).ToString( CultureInfo.InvariantCulture))); using var document = await SendAsync(uri, cancellationToken); var data = document.RootElement.GetProperty("data"); var result = data.GetProperty("result"); var samples = new List(); foreach (var item in result.EnumerateArray()) { var labels = ReadLabels(item.GetProperty("metric")); if (!TryReadValue(item.GetProperty("value"), out var value)) continue; samples.Add(new PrometheusSample(labels, value)); } return samples; } private async Task> QueryRangeAsync( string query, DateTime from, DateTime to, int stepSeconds, CancellationToken cancellationToken) { var uri = BuildUri( "api/v1/query_range", ("query", query), ("start", ToUnixSeconds(from).ToString( CultureInfo.InvariantCulture)), ("end", ToUnixSeconds(to).ToString( CultureInfo.InvariantCulture)), ("step", stepSeconds.ToString(CultureInfo.InvariantCulture))); using var document = await SendAsync(uri, cancellationToken); var result = document.RootElement .GetProperty("data") .GetProperty("result"); var first = result.EnumerateArray().FirstOrDefault(); if (first.ValueKind == JsonValueKind.Undefined || !first.TryGetProperty("values", out var values)) return []; var points = new List(); foreach (var value in values.EnumerateArray()) { if (!TryReadValue(value, out var measurement)) continue; var timestamp = value[0].GetDouble(); points.Add(new PerformanceSeriesPoint( DateTimeOffset.FromUnixTimeMilliseconds( checked((long)(timestamp * 1000))).UtcDateTime, measurement)); } return points; } private async Task SendAsync( Uri uri, CancellationToken cancellationToken) { using var request = new HttpRequestMessage(HttpMethod.Get, uri); if (!string.IsNullOrWhiteSpace(options.BearerToken)) { request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", options.BearerToken); } using var response = await httpClient.SendAsync( request, HttpCompletionOption.ResponseHeadersRead, cancellationToken); response.EnsureSuccessStatusCode(); await using var stream = await response.Content.ReadAsStreamAsync( cancellationToken); var document = await JsonDocument.ParseAsync( stream, cancellationToken: cancellationToken); if (!document.RootElement.TryGetProperty("status", out var status) || status.GetString() != "success") { document.Dispose(); throw new InvalidOperationException( "Prometheus 返回了非成功查询状态。"); } return document; } private Uri BuildUri( string relativePath, params (string Key, string Value)[] parameters) { var baseUri = new Uri( options.PrometheusBaseUrl.TrimEnd('/') + "/", UriKind.Absolute); var query = string.Join( "&", parameters.Select(parameter => $"{Uri.EscapeDataString(parameter.Key)}=" + $"{Uri.EscapeDataString(parameter.Value)}")); return new Uri(baseUri, $"{relativePath}?{query}"); } private static string Selector( string metric, params (string Label, string Operator, string Value)[] filters) { var matchers = string.Join( ",", filters.Select(filter => $"{filter.Label}{filter.Operator}\"" + $"{EscapePrometheusValue(filter.Value)}\"")); return $"{metric}{{{matchers}}}"; } private static string EscapePrometheusValue(string value) => value.Replace("\\", "\\\\", StringComparison.Ordinal) .Replace("\"", "\\\"", StringComparison.Ordinal) .Replace("\r", "\\r", StringComparison.Ordinal) .Replace("\n", "\\n", StringComparison.Ordinal); private static IReadOnlyDictionary ReadLabels( JsonElement metric) { var result = new Dictionary( StringComparer.Ordinal); foreach (var property in metric.EnumerateObject()) result[property.Name] = property.Value.GetString() ?? ""; return result; } private static bool TryReadValue( JsonElement value, out double measurement) { measurement = 0; if (value.ValueKind != JsonValueKind.Array || value.GetArrayLength() < 2) return false; var raw = value[1].GetString(); return double.TryParse( raw, NumberStyles.Float, CultureInfo.InvariantCulture, out measurement) && double.IsFinite(measurement); } private static IReadOnlyList MergeTimeline( IReadOnlyList requestRate, IReadOnlyList latency) { var points = new SortedDictionary(); foreach (var point in requestRate) { points[point.Timestamp] = new PerformanceTimelinePoint( point.Timestamp, Math.Round(point.Value, 3), null); } foreach (var point in latency) { points.TryGetValue(point.Timestamp, out var existing); points[point.Timestamp] = new PerformanceTimelinePoint( point.Timestamp, existing?.RequestsPerSecond, Math.Round(point.Value, 2)); } return points.Values.ToArray(); } private static IReadOnlyList MergeRanking( IReadOnlyList latency, IReadOnlyList count, IReadOnlyList exceptional, string label) { var names = latency .Concat(count) .Concat(exceptional) .Select(item => item.Labels.GetValueOrDefault(label)) .Where(name => !string.IsNullOrWhiteSpace(name)) .Distinct(StringComparer.Ordinal) .ToArray(); var items = names.Select(name => { var latencyValue = FindValue(latency, label, name); var countValue = FindValue(count, label, name); var exceptionalValue = FindValue(exceptional, label, name); return new PerformanceRankingItem( name!, Round(latencyValue), Round(countValue), Round(exceptionalValue)); }); return items .OrderByDescending(item => item.P95Milliseconds ?? -1) .ThenByDescending(item => item.RequestCount ?? -1) .Take(10) .ToArray(); } private static double? FindValue( IReadOnlyList samples, string label, string? name) => samples.FirstOrDefault(item => item.Labels.GetValueOrDefault(label) == name)?.Value; private static double ToUnixSeconds(DateTime value) => new DateTimeOffset( DateTime.SpecifyKind(value, DateTimeKind.Utc)).ToUnixTimeMilliseconds() / 1000d; private static double? Round(double? value) => value.HasValue && double.IsFinite(value.Value) ? Math.Round(value.Value, 2) : null; private static string? EmptyToNull(string? value) => string.IsNullOrWhiteSpace(value) ? null : value; private sealed record PrometheusSample( IReadOnlyDictionary Labels, double Value); private sealed record PerformanceSeriesPoint( DateTime Timestamp, double Value); } public sealed record PerformanceHeadline( double? RequestCount, double? RequestP95Milliseconds, double? ServerErrorRatePercent, double? DatabaseP95Milliseconds, double? SlowDatabaseCommandCount, double? FailedDatabaseCommandCount); public sealed record PerformanceTimelinePoint( DateTime Timestamp, double? RequestsPerSecond, double? RequestP95Milliseconds); public sealed record PerformanceRankingItem( string Name, double? P95Milliseconds, double? RequestCount, double? ExceptionalCount); public sealed record PerformanceReport( string Status, string Range, DateTime? From, DateTime? To, DateTime GeneratedAt, string DataSource, string? DashboardUrl, string? Detail, PerformanceHeadline? Headline, IReadOnlyList Timeline, IReadOnlyList Endpoints, IReadOnlyList DatabaseQueries) { public static PerformanceReport NotConfigured( string range, string? dashboardUrl) => Empty( "not_configured", range, dashboardUrl, "尚未配置 Prometheus 数据源。请先部署指标存储并设置 " + "PerformanceReporting__PrometheusBaseUrl。"); public static PerformanceReport Unavailable( string range, string? dashboardUrl) => Empty( "unavailable", range, dashboardUrl, "性能数据源暂时不可用。系统业务不受影响,请检查 Prometheus 与网络配置。"); private static PerformanceReport Empty( string status, string range, string? dashboardUrl, string detail) => new( status, range, null, null, DateTime.UtcNow, "prometheus", string.IsNullOrWhiteSpace(dashboardUrl) ? null : dashboardUrl, detail, null, [], [], []); } internal sealed record PerformanceRange( string Key, TimeSpan Duration, string PrometheusRange, string RateWindow, int StepSeconds) { public static PerformanceRange? TryParse(string? value) => value?.Trim().ToLowerInvariant() switch { "15m" => new("15m", TimeSpan.FromMinutes(15), "15m", "1m", 30), "1h" => new("1h", TimeSpan.FromHours(1), "1h", "5m", 60), "24h" => new("24h", TimeSpan.FromHours(24), "24h", "15m", 900), "7d" => new("7d", TimeSpan.FromDays(7), "7d", "1h", 3600), _ => null }; }