File: Telemetry\TelemetryShutdownTests.cs
Web Access
Project: src\tests\Aspire.Cli.Tests\Aspire.Cli.Tests.csproj (Aspire.Cli.Tests)
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
 
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Net;
using System.Text.Json;
using Aspire.Cli.Telemetry;
using Aspire.Cli.Tests.Utils;
using Aspire.Shared.Telemetry;
using Aspire.Tests.Shared.Telemetry;
using Azure.Core.Pipeline;
using Azure.Monitor.OpenTelemetry.Exporter;
using Microsoft.AspNetCore.InternalTesting;
using Microsoft.DotNet.RemoteExecutor;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Logging.Testing;
using OpenTelemetry;
using OpenTelemetry.Logs;
using OpenTelemetry.Trace;
 
namespace Aspire.Cli.Tests.Telemetry;
 
public class TelemetryShutdownTests(ITestOutputHelper outputHelper)
{
    [Theory]
    [InlineData(nameof(ReportedTelemetryMode.Local), CliExitCodes.Success, Timeout.Infinite)]
    [InlineData(nameof(ReportedTelemetryMode.Local), CliExitCodes.Cancelled, Timeout.Infinite)]
    [InlineData(nameof(ReportedTelemetryMode.Local), CliExitCodes.InvalidCommand, 300)]
    [InlineData(nameof(ReportedTelemetryMode.CI), CliExitCodes.Success, 5000)]
    [InlineData(nameof(ReportedTelemetryMode.CI), CliExitCodes.Cancelled, 5000)]
    [InlineData(nameof(ReportedTelemetryMode.CI), CliExitCodes.InvalidCommand, 5000)]
    [InlineData(nameof(ReportedTelemetryMode.Agent), CliExitCodes.Success, Timeout.Infinite)]
    [InlineData(nameof(ReportedTelemetryMode.Agent), CliExitCodes.InvalidCommand, Timeout.Infinite)]
    public void Shutdown_SelectsBudgetForBothSignalsAndRestoresDrainBudget(string modeName, int exitCode, int expectedTimeout)
    {
        using var process = RemoteExecutor.Invoke(static async (modeValue, exitValue, timeoutValue) =>
        {
            var mode = Enum.Parse<ReportedTelemetryMode>(modeValue);
            var exitCode = int.Parse(exitValue);
            var expectedTimeout = int.Parse(timeoutValue);
            TelemetryManager.ConfigureExporterForProcess(mode == ReportedTelemetryMode.Agent, mode == ReportedTelemetryMode.CI);
            var calls = new ConcurrentQueue<(string Signal, int Timeout, int DrainBudget)>();
            bool Shutdown(string signal, int timeout)
            {
                calls.Enqueue((signal, timeout,
                    Assert.IsType<int>(AppContext.GetData("Azure.Monitor.OpenTelemetry.Exporter.ShutdownDrainBudgetMilliseconds"))));
                return true;
            }
 
            var traceProcessor = new TestTelemetryProcessor<Activity> { ShutdownCallback = timeout => Shutdown("trace", timeout) };
            var logProcessor = new TestTelemetryProcessor<LogRecord> { ShutdownCallback = timeout => Shutdown("log", timeout) };
            using var fixture = new TelemetryFixture(initialize: false);
            using var manager = new TelemetryManager(
                new TelemetryConfiguration { ReportedTelemetryEnabled = true }, fixture.TagsSource, fixture.Telemetry,
                NullLogger<TelemetryManager>.Instance,
                (resource, _) => AzureMonitorTelemetryProvider.Create(new ServiceCollection(), resource, AspireCliTelemetry.EventLogCategoryName,
                    () => Sdk.CreateTracerProviderBuilder().AddProcessor(traceProcessor).Build(),
                    provider => provider.AddProcessor(logProcessor)));
            manager.Initialize();
 
            var shutdown = manager.TryShutdownAsync(exitCode, mode);
            Assert.Same(shutdown, manager.TryShutdownAsync(CliExitCodes.Success, mode));
            Assert.True(await shutdown);
            var expectedDrainBudget = expectedTimeout == 300 ? 300 : 0;
            Assert.Equal(
                [("log", expectedTimeout, expectedDrainBudget), ("trace", expectedTimeout, expectedDrainBudget)],
                calls.OrderBy(call => call.Signal).ToArray());
            Assert.Equal(0, AppContext.GetData("Azure.Monitor.OpenTelemetry.Exporter.ShutdownDrainBudgetMilliseconds"));
            manager.Dispose();
            Assert.Equal(1, traceProcessor.DisposeCount);
            Assert.Equal(1, logProcessor.DisposeCount);
        }, modeName, exitCode.ToString(), expectedTimeout.ToString());
    }
 
    [Fact]
    public void Shutdown_FailedDrainLogsWarningAndRestoresFailureBudget()
    {
        using var process = RemoteExecutor.Invoke(static async () =>
        {
            var mode = TelemetryManager.ConfigureExporterForProcess(isAgentTelemetryInvocation: false, isCIEnvironment: false);
            var sink = new TestSink();
            using var loggerFactory = LoggerFactory.Create(builder => builder.AddProvider(new TestLoggerProvider(sink)));
            using var fixture = new TelemetryFixture(initialize: false);
            var logProcessor = new TestTelemetryProcessor<LogRecord> { ShutdownCallback = static _ => false };
            using var manager = new TelemetryManager(
                new TelemetryConfiguration { ReportedTelemetryEnabled = true }, fixture.TagsSource, fixture.Telemetry,
                loggerFactory.CreateLogger<TelemetryManager>(),
                (resource, _) => AzureMonitorTelemetryProvider.Create(new ServiceCollection(), resource, AspireCliTelemetry.EventLogCategoryName,
                    () => Sdk.CreateTracerProviderBuilder().Build(), provider => provider.AddProcessor(logProcessor)));
            manager.Initialize();
 
            Assert.True(await manager.TryShutdownAsync(CliExitCodes.InvalidCommand, mode));
 
            var warning = Assert.Single(sink.Writes);
            Assert.Equal(LogLevel.Warning, warning.LogLevel);
            Assert.Equal("Timed out shutting down CLI reported telemetry.", warning.Message);
            Assert.Equal(0, AppContext.GetData("Azure.Monitor.OpenTelemetry.Exporter.ShutdownDrainBudgetMilliseconds"));
        });
    }
 
    [Theory]
    [InlineData("Local")]
    [InlineData("CI")]
    [InlineData("Agent")]
    [InlineData("Dashboard")]
    [InlineData("RetryableCI")]
    [InlineData("TimedOutCI")]
    [InlineData("UnavailableStorage")]
    public void Shutdown_HandlesBlockedDeliveryAndStorage(string scenario)
    {
        using var workspace = TemporaryWorkspace.CreateForCli(outputHelper);
        using var process = RemoteExecutor.Invoke(static async (directory, scenario) =>
        {
            Environment.SetEnvironmentVariable("APPLICATIONINSIGHTS_STATSBEAT_DISABLED", "true");
            var dashboard = scenario == "Dashboard";
            var isCI = scenario is "CI" or "RetryableCI" or "TimedOutCI";
            var mode = TelemetryManager.ConfigureExporterForProcess(scenario == "Agent", isCI);
            if (dashboard)
            {
                // The dashboard deliberately retains the service defaults, not the CLI's zero drain.
                AppContext.SetSwitch("Azure.Monitor.OpenTelemetry.Exporter.PersistOnForceFlush", false);
                AppContext.SetData("Azure.Monitor.OpenTelemetry.Exporter.ShutdownDrainBudgetMilliseconds", null);
            }
            var storage = Path.Combine(directory, "storage");
            if (scenario == "UnavailableStorage")
            {
                await File.WriteAllTextAsync(storage, "A file occupies the requested storage directory.");
            }
            var requestStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
            var releaseResponse = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
            var delivered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
            var exported = new ConcurrentQueue<(string Type, string Name)>();
            using var handler = new MockHttpMessageHandler(async (request, cancellationToken) =>
            {
                var payload = await request.Content!.ReadAsStringAsync(cancellationToken);
                // Azure sends newline-delimited envelopes:
                // {"data":{"baseType":"MessageData","baseData":{"message":"event-0",...}},...}
                // Trace envelopes use RemoteDependencyData/name; resource metrics share this transport.
                var lines = payload.Split('\n', StringSplitOptions.RemoveEmptyEntries);
                var records = new List<(string Type, string Name)>();
                foreach (var line in lines)
                {
                    using var envelope = JsonDocument.Parse(line);
                    var data = envelope.RootElement.GetProperty("data");
                    var type = data.GetProperty("baseType").GetString();
                    var baseData = data.GetProperty("baseData");
                    if (type is "MessageData" or "RemoteDependencyData")
                    {
                        records.Add((type, baseData.GetProperty(type == "MessageData" ? "message" : "name").GetString()!));
                    }
                }
                if (records.Count > 0)
                {
                    requestStarted.TrySetResult();
                    await releaseResponse.Task.WaitAsync(cancellationToken);
                    if (scenario == "RetryableCI")
                    {
                        return new HttpResponseMessage(HttpStatusCode.ServiceUnavailable);
                    }
                    foreach (var record in records)
                    {
                        exported.Enqueue(record);
                    }
                    if (exported.Count == 6)
                    {
                        delivered.TrySetResult();
                    }
                }
                return new HttpResponseMessage(HttpStatusCode.OK)
                {
                    Content = new StringContent(JsonSerializer.Serialize(new { itemsReceived = lines.Length, itemsAccepted = lines.Length, errors = Array.Empty<object>() }))
                };
            });
            using var client = new HttpClient(handler);
            void ConfigureTransport(AzureMonitorExporterOptions options)
            {
                options.Transport = new HttpClientTransport(client);
                // Exercise the exporter's retryable-response storage behavior without Azure.Core backoff.
                options.Retry.MaxRetries = 0;
            }
 
            using var fixture = new TelemetryFixture(initialize: false);
            AzureMonitorTelemetryProvider? reportedProvider = null;
            AzureMonitorExporterOptions? traceOptions = null;
            using var manager = new TelemetryManager(
                new TelemetryConfiguration { ReportedTelemetryEnabled = true }, fixture.TagsSource, fixture.Telemetry,
                NullLogger<TelemetryManager>.Instance,
                (resource, _) => reportedProvider = AzureMonitorTelemetryProvider.Create(new ServiceCollection(), resource,
                    AspireCliTelemetry.ReportedActivitySourceName, AspireCliTelemetry.EventLogCategoryName,
                    $"InstrumentationKey={Guid.NewGuid()};IngestionEndpoint=https://localhost/", storage,
                    builder => builder.ConfigureServices(services => services.Configure<AzureMonitorExporterOptions>(options =>
                    {
                        traceOptions = options;
                        ConfigureTransport(options);
                    })),
                    options =>
                    {
                        Assert.Equal(TimeSpan.FromSeconds(5), options.Retry.NetworkTimeout);
                        ConfigureTransport(options);
                    }));
            manager.Initialize();
            Assert.NotNull(reportedProvider);
            Assert.NotNull(traceOptions);
            Assert.Equal(TimeSpan.FromSeconds(5), traceOptions.Retry.NetworkTimeout);
            using var source = new ActivitySource(AspireCliTelemetry.ReportedActivitySourceName);
            for (var i = 0; i < 3; i++)
            {
                using var activity = source.StartActivity($"operation-{i}");
                Assert.NotNull(activity);
                fixture.Telemetry.RecordEvent($"event-{i}");
            }
 
            try
            {
                var shutdown = dashboard ? reportedProvider.ShutdownAsync(5000)
                    : manager.TryShutdownAsync(CliExitCodes.Success, mode);
                await requestStarted.Task.DefaultTimeout();
                if (scenario is "Local" or "Agent" or "Dashboard")
                {
                    Assert.True(await shutdown.DefaultTimeout());
                    Assert.False(delivered.Task.IsCompleted);
                    Assert.NotEmpty(Directory.EnumerateFiles(storage, "*", SearchOption.AllDirectories));
                }
                else if (scenario == "TimedOutCI")
                {
                    // Provider and network budgets must permit shutdown without a response. Its result
                    // means lifecycle completion, not ingestion acceptance or a changed command result.
                    Assert.True(await shutdown.DefaultTimeout());
                    Assert.Empty(exported);
                    return;
                }
                else
                {
                    Assert.False(shutdown.IsCompleted);
                }
 
                releaseResponse.SetResult();
                Assert.True(await shutdown.DefaultTimeout());
                if (scenario == "RetryableCI")
                {
                    Assert.Empty(exported);
                    Assert.NotEmpty(Directory.EnumerateFiles(storage, "*", SearchOption.AllDirectories));
                }
                else
                {
                    await delivered.Task.DefaultTimeout();
                    Assert.Equal(
                        [("MessageData", "event-0"), ("MessageData", "event-1"), ("MessageData", "event-2"),
                         ("RemoteDependencyData", "operation-0"), ("RemoteDependencyData", "operation-1"), ("RemoteDependencyData", "operation-2")],
                        exported.OrderBy(record => record.Type).ThenBy(record => record.Name).ToArray());
                }
            }
            finally
            {
                releaseResponse.TrySetResult();
                await manager.TryShutdownAsync(CliExitCodes.Success, mode).DefaultTimeout();
            }
        }, workspace.Path, scenario);
    }
}