File: Otlp\Storage\SqliteTelemetryRepository.Metrics.Reads.cs
Web Access
Project: src\src\Aspire.Dashboard\Aspire.Dashboard.csproj (Aspire.Dashboard)
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
 
using System.Globalization;
using System.Text;
using Aspire.Dashboard.Otlp.Model;
using Aspire.Dashboard.Otlp.Model.MetricValues;
using Dapper;
using Microsoft.Data.Sqlite;
 
namespace Aspire.Dashboard.Otlp.Storage;
 
public sealed partial class SqliteTelemetryRepository
{
    private const int MaxMetricReadBatchSize = 500;
    private const string MetricPointRangeFilterSql = "p.start_time_ticks <= @EndTicks AND p.end_time_ticks >= r.start_time_ticks";
    private const string MetricSourcePointRangeFilterSql = "source.start_time_ticks <= @EndTicks AND source.end_time_ticks >= r.start_time_ticks";
    private const string SelectedMetricPointsCteSql = $"""
        selected_metric_points AS (
            SELECT
                p.point_id,
                p.dimension_id,
                p.point_type,
                p.start_time_ticks,
                p.end_time_ticks,
                p.repeat_count,
                p.integer_value,
                p.double_value,
                p.histogram_sum,
                p.histogram_count
            FROM telemetry_metric_points p
            JOIN metric_dimension_query_ranges r ON r.dimension_id = p.dimension_id
            WHERE {MetricPointRangeFilterSql}
        )
        """;
    private const string EffectiveMetricPointsCteSql = $"""
        {SelectedMetricPointsCteSql},
        ranked_metric_points AS (
            SELECT
                p.*,
                ROW_NUMBER() OVER (
                    PARTITION BY p.start_time_ticks, p.dimension_id
                    ORDER BY p.point_id DESC) AS point_rank
            FROM selected_metric_points p
        ),
        effective_metric_points AS (
            SELECT *
            FROM ranked_metric_points
            WHERE point_rank = 1
        )
        """;
    private static readonly string s_rolledUpMetricPointsCteSql = $"""
        {EffectiveMetricPointsCteSql},
        bucketed_metric_points AS (
            SELECT
                p.*,
                (p.start_time_ticks / @PointIntervalTicks) * @PointIntervalTicks AS rollup_start_time_ticks
            FROM effective_metric_points p
        ),
        ranked_rollup_metric_points AS (
            SELECT
                p.*,
                MAX(p.end_time_ticks) OVER (
                    PARTITION BY p.dimension_id, p.point_type, p.rollup_start_time_ticks) AS rollup_end_time_ticks,
                SUM(p.repeat_count) OVER (
                    PARTITION BY p.dimension_id, p.point_type, p.rollup_start_time_ticks) AS rollup_repeat_count,
                ROW_NUMBER() OVER (
                    PARTITION BY p.dimension_id, p.point_type, p.rollup_start_time_ticks
                    ORDER BY
                        CASE WHEN p.point_type = {HistogramPointType} THEN p.start_time_ticks END DESC,
                        p.integer_value DESC,
                        p.double_value DESC,
                        p.point_id DESC) AS rollup_rank
            FROM bucketed_metric_points p
        ),
        rolled_up_metric_points AS (
            SELECT
                p.point_id,
                p.dimension_id,
                p.point_type,
                p.rollup_start_time_ticks AS start_time_ticks,
                p.rollup_end_time_ticks AS end_time_ticks,
                p.rollup_repeat_count AS repeat_count,
                p.integer_value,
                p.double_value,
                p.histogram_sum,
                p.histogram_count
            FROM ranked_rollup_metric_points p
            WHERE p.rollup_rank = 1
        )
        """;
    private const string FullFidelityMetricPointsCteSql = $"""
        {EffectiveMetricPointsCteSql},
        rolled_up_metric_points AS (
            SELECT
                p.point_id,
                p.dimension_id,
                p.point_type,
                p.start_time_ticks,
                p.end_time_ticks,
                p.repeat_count,
                p.integer_value,
                p.double_value,
                p.histogram_sum,
                p.histogram_count
            FROM effective_metric_points p
        )
        """;
 
    private OtlpInstrumentData? GetInstrumentFromDatabase(GetInstrumentRequest request, CancellationToken cancellationToken)
    {
        var instruments = GetCachedInstruments(request.ResourceKey, request.MeterName, request.InstrumentName);
        if (instruments.Count == 0)
        {
            return null;
        }
 
        using var connection = _database.OpenConnection();
        using var interrupt = connection.RegisterInterrupt(cancellationToken);
        var knownAttributeValues = new Dictionary<string, List<string?>>();
        var dimensions = MaterializeMetricDimensions(
            connection,
            instruments,
            request.StartTime,
            request.EndTime,
            request.DimensionFilters,
            request.DimensionCursors,
            request.DataPointInterval,
            request.IncludeExemplars,
            request.PopulateExemplarAttributes,
            knownAttributeValues,
            out var hasOverflow);
        return new OtlpInstrumentData
        {
            Summary = instruments[0].Summary,
            Dimensions = dimensions,
            KnownAttributeValues = knownAttributeValues,
            HasOverflow = hasOverflow
        };
    }
 
    private DateTime? GetInstrumentLatestEndTimeFromDatabase(ResourceKey resourceKey, string meterName, string instrumentName)
    {
        using var connection = _database.OpenConnection();
        var endTimeTicks = connection.QuerySingleOrDefault<long?>("""
            SELECT MAX(p.end_time_ticks)
            FROM telemetry_metric_points p
            JOIN telemetry_metric_dimensions d ON d.dimension_id = p.dimension_id
            JOIN telemetry_metric_instruments i ON i.instrument_id = d.instrument_id
            JOIN telemetry_resources r ON r.resource_id = i.resource_id
            JOIN telemetry_scopes s ON s.scope_id = i.scope_id
            WHERE r.resource_name = @ResourceName COLLATE NOCASE
                            AND (@InstanceId IS NULL OR r.instance_id = @InstanceId COLLATE NOCASE)
              AND s.scope_name = @MeterName
              AND i.instrument_name = @InstrumentName;
            """, new { ResourceName = resourceKey.Name, resourceKey.InstanceId, MeterName = meterName, InstrumentName = instrumentName });
        return endTimeTicks is not null ? new DateTime(endTimeTicks.Value, DateTimeKind.Utc) : null;
    }
 
    private List<DimensionScope> MaterializeMetricDimensions(
        SqliteConnection connection,
        IReadOnlyList<CachedInstrument> instruments,
        DateTime? startTime,
        DateTime? endTime,
        IReadOnlyDictionary<string, IReadOnlyList<string?>> dimensionFilters,
        IReadOnlyList<MetricDimensionCursor> dimensionCursors,
        TimeSpan? dataPointInterval,
        bool includeExemplars,
        bool populateExemplarAttributes,
        Dictionary<string, List<string?>> knownAttributeValues,
        out bool hasOverflow)
    {
        if (dataPointInterval is { } interval && interval <= TimeSpan.Zero)
        {
            throw new ArgumentOutOfRangeException(nameof(dataPointInterval), interval, "The metric data point interval must be greater than zero.");
        }
 
        var dimensions = GetMetricDimensions(
            connection,
            instruments,
            dimensionFilters,
            knownAttributeValues,
            out hasOverflow);
        if (dimensions.Count == 0)
        {
            return [];
        }
 
        var results = dimensions
            .Select(dimension => dimension.Scope)
            .ToList();
        if (startTime is null || endTime is null)
        {
            return results;
        }
 
        var queryParameters = new DynamicParameters();
        queryParameters.Add("EndTicks", endTime.Value.Ticks);
        queryParameters.Add("PointIntervalTicks", dataPointInterval?.Ticks ?? 0);
        var dimensionQueryRangesCteSql = CreateMetricDimensionQueryRangesCte(
            dimensions,
            dimensionCursors,
            startTime,
            queryParameters);
        var metricPointsCteSql = dataPointInterval is null ? FullFidelityMetricPointsCteSql : s_rolledUpMetricPointsCteSql;
        var pointRecords = connection.Query<MetricPointDataRecord>($"""
            WITH {dimensionQueryRangesCteSql},
            {metricPointsCteSql}
            SELECT
                p.point_id AS PointId,
                p.dimension_id AS DimensionId,
                p.point_type AS PointType,
                p.start_time_ticks AS StartTimeTicks,
                p.end_time_ticks AS EndTimeTicks,
                p.repeat_count AS RepeatCount,
                p.integer_value AS IntegerValue,
                p.double_value AS DoubleValue,
                p.histogram_sum AS HistogramSum,
                p.histogram_count AS HistogramCount,
                stored.bucket_counts AS BucketCounts,
                stored.explicit_bounds AS ExplicitBounds
            FROM rolled_up_metric_points p
            JOIN telemetry_metric_points stored ON stored.point_id = p.point_id
            ORDER BY p.dimension_id, p.start_time_ticks, p.point_id;
            """, queryParameters).AsList();
        var points = pointRecords.ToLookup(record => record.DimensionId);
        var exemplars = includeExemplars
            ? MaterializeMetricExemplars(
                connection,
                dimensionQueryRangesCteSql,
                queryParameters,
                dataPointInterval,
                pointRecords,
                populateExemplarAttributes)
            : Array.Empty<KeyValuePair<long, MetricsExemplar>>().ToLookup(pair => pair.Key, pair => pair.Value);
 
        for (var dimensionIndex = 0; dimensionIndex < dimensions.Count; dimensionIndex++)
        {
            var dimensionId = dimensions[dimensionIndex].DimensionId;
            var dimension = results[dimensionIndex];
            foreach (var point in points[dimensionId])
            {
                MetricValueBase value = point.PointType switch
                {
                    LongPointType => new MetricValue<long>(point.IntegerValue!.Value, new DateTime(point.StartTimeTicks, DateTimeKind.Utc), new DateTime(point.EndTimeTicks, DateTimeKind.Utc)),
                    DoublePointType => new MetricValue<double>(point.DoubleValue!.Value, new DateTime(point.StartTimeTicks, DateTimeKind.Utc), new DateTime(point.EndTimeTicks, DateTimeKind.Utc)),
                    HistogramPointType => CreateHistogramValue(point),
                    _ => throw new InvalidOperationException($"Unknown metric point type '{point.PointType}'.")
                };
                if (point.PointType != HistogramPointType)
                {
                    value.Count = checked((ulong)point.RepeatCount);
                }
                value.Exemplars.AddRange(exemplars[point.PointId]);
                dimension.Values.Add(value);
            }
        }
        return results;
    }
 
    private List<StoredMetricDimension> GetMetricDimensions(
        SqliteConnection connection,
        IReadOnlyList<CachedInstrument> instruments,
        IReadOnlyDictionary<string, IReadOnlyList<string?>> dimensionFilters,
        Dictionary<string, List<string?>> knownAttributeValues,
        out bool hasOverflow)
    {
        var dimensionRecords = connection.Query<MetricDimensionAttributeRecord>("""
            SELECT
                d.dimension_id AS DimensionId,
                d.instrument_id AS InstrumentId,
                a.attribute_key AS AttributeKey,
                a.attribute_value AS AttributeValue
            FROM telemetry_metric_dimensions d
            LEFT JOIN telemetry_metric_dimension_attributes a ON a.dimension_id = d.dimension_id
            WHERE d.instrument_id IN @InstrumentIds
            ORDER BY d.dimension_id, a.ordinal;
            """, new { InstrumentIds = instruments.Select(instrument => instrument.InstrumentId) }).AsList();
        var dimensionIds = dimensionRecords.Select(record => record.DimensionId).Distinct().ToArray();
        var pointAttributes = dimensionRecords
            .Where(record => record.AttributeKey is not null)
            .Select(record => new OwnedAttributeRecord
            {
                OwnerId = record.DimensionId,
                AttributeKey = record.AttributeKey!,
                AttributeValue = record.AttributeValue!
            })
            .ToLookup(record => record.OwnerId);
        hasOverflow = dimensionIds.Any(dimensionId => IsOverflowDimension(pointAttributes[dimensionId]));
        var instrumentsById = instruments.ToDictionary(instrument => instrument.InstrumentId);
        var instrumentIdsByDimensionId = dimensionRecords
            .DistinctBy(record => record.DimensionId)
            .ToDictionary(record => record.DimensionId, record => record.InstrumentId);
        // Dimension rows store only point attributes. Scope attributes are stored once with the scope and merged here for display and filtering.
        var attributes = dimensionIds
            .SelectMany(dimensionId => pointAttributes[dimensionId]
                .Select(attribute => KeyValuePair.Create(attribute.AttributeKey, attribute.AttributeValue))
                .Concat(instrumentsById[instrumentIdsByDimensionId[dimensionId]].Summary.Parent.Attributes)
                .Select(attribute => new OwnedAttributeRecord
                {
                    OwnerId = dimensionId,
                    AttributeKey = attribute.Key,
                    AttributeValue = attribute.Value
                }))
            .ToLookup(record => record.OwnerId);
        PopulateKnownAttributeValues(dimensionIds, attributes, knownAttributeValues);
 
        return dimensionIds
            .Where(dimensionId => MatchesDimensionFilters(attributes[dimensionId], dimensionFilters))
            .Select(dimensionId => new StoredMetricDimension(
                dimensionId,
                new DimensionScope(
                    _otlpContext.Options.MaxMetricsCount,
                    attributes[dimensionId].Select(attribute => KeyValuePair.Create(attribute.AttributeKey, attribute.AttributeValue)).ToArray())))
            .ToList();
    }
 
    private static bool IsOverflowDimension(IEnumerable<OwnedAttributeRecord> attributes)
    {
        using var enumerator = attributes.GetEnumerator();
        return enumerator.MoveNext() &&
            enumerator.Current is { AttributeKey: "otel.metric.overflow", AttributeValue: "true" } &&
            !enumerator.MoveNext();
    }
 
    private static string CreateMetricDimensionQueryRangesCte(
        IReadOnlyList<StoredMetricDimension> dimensions,
        IReadOnlyList<MetricDimensionCursor> dimensionCursors,
        DateTime? defaultStartTime,
        DynamicParameters queryParameters)
    {
        var sql = new StringBuilder("metric_dimension_query_ranges(dimension_id, start_time_ticks) AS (VALUES ");
        for (var i = 0; i < dimensions.Count; i++)
        {
            if (i > 0)
            {
                sql.Append(", ");
            }
 
            var dimension = dimensions[i];
            var cursor = dimensionCursors.FirstOrDefault(cursor =>
                cursor.Attributes.SequenceEqual(dimension.Scope.Attributes));
            queryParameters.Add($"DimensionId{i}", dimension.DimensionId);
            queryParameters.Add($"DimensionStartTicks{i}", cursor?.StartTime.Ticks ?? defaultStartTime?.Ticks ?? 0);
            sql.Append(CultureInfo.InvariantCulture, $"(@DimensionId{i}, @DimensionStartTicks{i})");
        }
        sql.Append(')');
        return sql.ToString();
    }
 
    private static bool MatchesDimensionFilters(IEnumerable<OwnedAttributeRecord> attributes, IReadOnlyDictionary<string, IReadOnlyList<string?>> dimensionFilters)
    {
        foreach (var (key, values) in dimensionFilters)
        {
            var value = attributes.FirstOrDefault(attribute => attribute.AttributeKey == key)?.AttributeValue;
            if (!values.Contains(value))
            {
                return false;
            }
        }
        return true;
    }
 
    private static void PopulateKnownAttributeValues(
        IReadOnlyList<long> dimensionIds,
        ILookup<long, OwnedAttributeRecord> attributes,
        Dictionary<string, List<string?>> knownAttributeValues)
    {
        // Point and scope attributes were already accepted during ingestion, so intentionally do not limit their
        // merged key or per-key value counts while building display metadata.
        for (var dimensionIndex = 0; dimensionIndex < dimensionIds.Count; dimensionIndex++)
        {
            var dimensionId = dimensionIds[dimensionIndex];
            foreach (var key in knownAttributeValues.Keys.Union(attributes[dimensionId].Select(attribute => attribute.AttributeKey)).Distinct().ToList())
            {
                if (!knownAttributeValues.TryGetValue(key, out var values))
                {
                    values = [];
                    knownAttributeValues.Add(key, values);
                    if (dimensionIndex > 0)
                    {
                        TryAddValue(values, null);
                    }
                }
                var value = attributes[dimensionId].FirstOrDefault(attribute => attribute.AttributeKey == key)?.AttributeValue;
                TryAddValue(values, value);
            }
        }
 
        static void TryAddValue(List<string?> values, string? value)
        {
            if (!values.Contains(value))
            {
                values.Add(value);
            }
        }
    }
 
    private static HistogramValue CreateHistogramValue(MetricPointDataRecord point) => new(
        UnpackUInt64Values(point.BucketCounts!),
        point.HistogramSum!.Value,
        checked((ulong)point.HistogramCount!.Value),
        new DateTime(point.StartTimeTicks, DateTimeKind.Utc),
        new DateTime(point.EndTimeTicks, DateTimeKind.Utc),
        UnpackDoubleValues(point.ExplicitBounds!));
 
    private static ILookup<long, MetricsExemplar> MaterializeMetricExemplars(
        SqliteConnection connection,
        string dimensionQueryRangesCteSql,
        DynamicParameters queryParameters,
        TimeSpan? dataPointInterval,
        IReadOnlyList<MetricPointDataRecord> points,
        bool populateExemplarAttributes)
    {
        var records = connection.Query<MetricExemplarRecord>($"""
            WITH {dimensionQueryRangesCteSql}
            SELECT
                e.exemplar_id AS ExemplarId,
                source.dimension_id AS DimensionId,
                source.point_type AS PointType,
                source.start_time_ticks AS SourceStartTimeTicks,
                e.start_time_ticks AS StartTimeTicks,
                e.exemplar_value AS ExemplarValue,
                e.span_id AS SpanId,
                e.trace_id AS TraceId
            FROM telemetry_metric_exemplars e
            JOIN telemetry_metric_points source ON source.point_id = e.point_id
            JOIN metric_dimension_query_ranges r ON r.dimension_id = source.dimension_id
            WHERE {MetricSourcePointRangeFilterSql}
              AND e.start_time_ticks >= r.start_time_ticks
              AND e.start_time_ticks <= @EndTicks
            ORDER BY source.dimension_id, source.start_time_ticks, e.exemplar_id;
            """, queryParameters).AsList();
        var pointIds = points.ToDictionary(
            point => new MetricPointKey(point.DimensionId, point.PointType, point.StartTimeTicks),
            point => point.PointId);
        var pointIntervalTicks = dataPointInterval?.Ticks;
        var mappedRecords = new List<(long PointId, MetricExemplarRecord Record)>();
        foreach (var record in records)
        {
            var rollupStartTimeTicks = pointIntervalTicks is { } intervalTicks
                ? (record.SourceStartTimeTicks / intervalTicks) * intervalTicks
                : record.SourceStartTimeTicks;
            if (pointIds.TryGetValue(new MetricPointKey(record.DimensionId, record.PointType, rollupStartTimeTicks), out var pointId))
            {
                mappedRecords.Add((pointId, record));
            }
        }
        var attributes = populateExemplarAttributes
            ? MaterializeMetricExemplarAttributes(connection, mappedRecords.Select(item => item.Record.ExemplarId).Distinct().ToArray())
            : Array.Empty<OwnedAttributeRecord>().ToLookup(record => record.OwnerId);
        return mappedRecords
            .OrderBy(item => item.PointId)
            .ThenBy(item => item.Record.ExemplarId)
            .Select(item => new KeyValuePair<long, MetricsExemplar>(
                item.PointId,
                new MetricsExemplar
                {
                    Start = new DateTime(item.Record.StartTimeTicks, DateTimeKind.Utc),
                    Value = item.Record.ExemplarValue,
                    SpanId = item.Record.SpanId,
                    TraceId = item.Record.TraceId,
                    Attributes = attributes[item.Record.ExemplarId].Select(attribute => KeyValuePair.Create(attribute.AttributeKey, attribute.AttributeValue)).ToArray()
                }))
            .ToLookup(pair => pair.Key, pair => pair.Value);
    }
 
    private static ILookup<long, OwnedAttributeRecord> MaterializeMetricExemplarAttributes(
        SqliteConnection connection,
        IReadOnlyList<long> exemplarIds)
    {
        var attributes = new List<OwnedAttributeRecord>();
        foreach (var exemplarIdBatch in exemplarIds.Chunk(MaxMetricReadBatchSize))
        {
            attributes.AddRange(connection.Query<OwnedAttributeRecord>("""
                SELECT exemplar_id AS OwnerId, attribute_key AS AttributeKey, attribute_value AS AttributeValue
                FROM telemetry_metric_exemplar_attributes
                WHERE exemplar_id IN @ExemplarIds
                ORDER BY exemplar_id, ordinal;
                """, new { ExemplarIds = exemplarIdBatch }));
        }
        return attributes.ToLookup(record => record.OwnerId);
    }
 
    private sealed class MetricPointDataRecord : MetricPointRecord
    {
        public required long DimensionId { get; init; }
        public required long StartTimeTicks { get; init; }
        public required long RepeatCount { get; init; }
        public double? HistogramSum { get; init; }
        public byte[]? BucketCounts { get; init; }
        public byte[]? ExplicitBounds { get; init; }
    }
 
    private sealed class MetricDimensionAttributeRecord
    {
        public required long DimensionId { get; init; }
        public required long InstrumentId { get; init; }
        public string? AttributeKey { get; init; }
        public string? AttributeValue { get; init; }
    }
 
    private sealed class MetricExemplarRecord
    {
        public required long ExemplarId { get; init; }
        public required long DimensionId { get; init; }
        public required int PointType { get; init; }
        public required long SourceStartTimeTicks { get; init; }
        public required long StartTimeTicks { get; init; }
        public required double ExemplarValue { get; init; }
        public required string SpanId { get; init; }
        public required string TraceId { get; init; }
    }
 
    private sealed record StoredMetricDimension(long DimensionId, DimensionScope Scope);
 
    private readonly record struct MetricPointKey(long DimensionId, int PointType, long StartTimeTicks);
}