File: Commands\TelemetryFollowTests.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.Net;
using System.Net.Http.Headers;
using System.Net.Sockets;
using System.Text;
using Aspire.Cli.Backchannel;
using Aspire.Cli.Commands;
using Aspire.Cli.Interaction;
using Aspire.Cli.Resources;
using Aspire.Cli.Tests.TestServices;
using Aspire.Cli.Tests.Utils;
using Microsoft.AspNetCore.InternalTesting;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
 
namespace Aspire.Cli.Tests.Commands;
 
public class TelemetryFollowTests(ITestOutputHelper outputHelper)
{
    [Theory]
    [InlineData("logs", false, false)]
    [InlineData("logs", false, true)]
    [InlineData("logs", true, false)]
    [InlineData("logs", true, true)]
    [InlineData("spans", false, false)]
    [InlineData("spans", false, true)]
    [InlineData("spans", true, false)]
    [InlineData("spans", true, true)]
    public async Task Follow_ExpectedDisconnect_ReturnsSuccessAndWritesStatusToStderr(string commandName, bool standalone, bool connectionReset)
    {
        var exception = connectionReset
            ? new IOException("Connection reset.", new SocketException((int)SocketError.ConnectionReset))
            : new HttpIOException(HttpRequestError.ResponseEnded, "The response ended prematurely.");
 
        var (exitCode, interactionService) = await InvokeAsync(commandName, standalone, follow: true, exception);
 
        Assert.Equal(CliExitCodes.Success, exitCode);
        Assert.Empty(interactionService.DisplayedErrors);
        Assert.Collection(interactionService.DisplayedRawText,
            output =>
            {
                Assert.Equal("[]", output.Text);
                Assert.Equal(ConsoleOutput.Standard, output.ConsoleOverride);
            },
            output =>
            {
                Assert.Equal(TelemetryCommandStrings.DashboardConnectionLost, output.Text);
                Assert.Equal(ConsoleOutput.Error, output.ConsoleOverride);
            });
    }
 
    [Theory]
    [InlineData("logs", true)]
    [InlineData("logs", false)]
    [InlineData("spans", true)]
    [InlineData("spans", false)]
    public async Task Follow_UnexpectedReadFailure_ReturnsFailure(string commandName, bool invalidResponse)
    {
        IOException exception = invalidResponse
            ? new HttpIOException(HttpRequestError.InvalidResponse, "Invalid response.")
            : new IOException("Unexpected read failure.");
 
        var (exitCode, interactionService) = await InvokeAsync(commandName, standalone: true, follow: true, exception);
 
        Assert.Equal(CliExitCodes.InvalidCommand, exitCode);
        Assert.Single(interactionService.DisplayedErrors);
        Assert.Collection(interactionService.DisplayedRawText,
            output =>
            {
                Assert.Equal("[]", output.Text);
                Assert.Equal(ConsoleOutput.Standard, output.ConsoleOverride);
            });
    }
 
    [Theory]
    [InlineData("logs")]
    [InlineData("spans")]
    public async Task Snapshot_TruncatedResponse_ReturnsFailure(string commandName)
    {
        var exception = new HttpIOException(HttpRequestError.ResponseEnded, "The response ended prematurely.");
 
        var (exitCode, interactionService) = await InvokeAsync(commandName, standalone: true, follow: false, exception);
 
        Assert.NotEqual(CliExitCodes.Success, exitCode);
        Assert.Single(interactionService.DisplayedErrors);
        Assert.Empty(interactionService.DisplayedRawText);
    }
 
    [Theory]
    [InlineData("logs")]
    [InlineData("spans")]
    public async Task Follow_InvalidTelemetry_ReturnsFailure(string commandName)
    {
        var exception = new HttpIOException(HttpRequestError.ResponseEnded, "The response ended prematurely.");
 
        var (exitCode, interactionService) = await InvokeAsync(commandName, standalone: true, follow: true, exception, "invalid json\n");
 
        Assert.Equal(CliExitCodes.InvalidCommand, exitCode);
        Assert.Single(interactionService.DisplayedErrors);
        Assert.Empty(interactionService.DisplayedRawText);
    }
 
    [Theory]
    [InlineData(false)]
    [InlineData(true)]
    public async Task ReadLinesWithDisconnectHandling_Cancellation_Propagates(bool disconnect)
    {
        using var cts = new CancellationTokenSource();
        Exception exception = disconnect
            ? new HttpIOException(HttpRequestError.ResponseEnded, "The response ended prematurely.")
            : new OperationCanceledException(cts.Token);
        using var stream = new FaultingReadStream([], exception);
        using var reader = new StreamReader(stream);
        var interactionService = new TestInteractionService();
        await using var enumerator = reader.ReadLinesWithDisconnectHandlingAsync(interactionService, cts.Token).GetAsyncEnumerator();
 
        if (disconnect)
        {
            cts.Cancel();
            Assert.False(await enumerator.MoveNextAsync());
        }
        else
        {
            await Assert.ThrowsAsync<OperationCanceledException>(async () => await enumerator.MoveNextAsync());
        }
 
        Assert.Empty(interactionService.DisplayedRawText);
    }
 
    [Fact]
    public async Task ReadLinesWithDisconnectHandling_CompleteResponse_YieldsLinesWithoutDisconnectStatus()
    {
        using var stream = new MemoryStream(Encoding.UTF8.GetBytes("first\n\nsecond\n"));
        using var reader = new StreamReader(stream);
        var interactionService = new TestInteractionService();
        var lines = new List<string>();
 
        await foreach (var line in reader.ReadLinesWithDisconnectHandlingAsync(interactionService, CancellationToken.None))
        {
            lines.Add(line);
        }
 
        Assert.Equal(["first", "second"], lines);
        Assert.Empty(interactionService.DisplayedRawText);
    }
 
    [Theory]
    [InlineData(false)]
    [InlineData(true)]
    public async Task ReadLines_WithoutErrorHandler_PropagatesReadFailure(bool explicitNull)
    {
        var exception = new HttpIOException(HttpRequestError.ResponseEnded, "The response ended prematurely.");
        using var stream = new FaultingReadStream([], exception);
        using var reader = new StreamReader(stream);
        var lines = explicitNull
            ? reader.ReadLinesAsync(onError: null, CancellationToken.None)
            : reader.ReadLinesAsync(CancellationToken.None);
        await using var enumerator = lines.GetAsyncEnumerator();
 
        var thrown = await Assert.ThrowsAsync<HttpIOException>(async () => await enumerator.MoveNextAsync());
 
        Assert.Same(exception, thrown);
    }
 
    [Theory]
    [InlineData(false)]
    [InlineData(true)]
    public async Task ReadLines_ErrorHandler_ControlsReadFailureHandling(bool handled)
    {
        var exception = new IOException("Read failure.");
        using var stream = new FaultingReadStream(Encoding.UTF8.GetBytes("first\n\nsecond\n"), exception);
        using var reader = new StreamReader(stream);
        var errors = new List<Exception>();
        await using var enumerator = reader.ReadLinesAsync(ex =>
        {
            errors.Add(ex);
            return handled;
        }, CancellationToken.None).GetAsyncEnumerator();
 
        Assert.True(await enumerator.MoveNextAsync());
        Assert.Equal("first", enumerator.Current);
        Assert.True(await enumerator.MoveNextAsync());
        Assert.Equal("second", enumerator.Current);
        Assert.Empty(errors);
 
        if (handled)
        {
            Assert.False(await enumerator.MoveNextAsync());
            Assert.False(await enumerator.MoveNextAsync());
        }
        else
        {
            var thrown = await Assert.ThrowsAsync<IOException>(async () => await enumerator.MoveNextAsync());
            Assert.Same(exception, thrown);
        }
 
        Assert.Collection(errors, error => Assert.Same(exception, error));
    }
 
    [Fact]
    public async Task ReadLines_ErrorHandlerFailure_Propagates()
    {
        using var stream = new FaultingReadStream([], new IOException("Read failure."));
        using var reader = new StreamReader(stream);
        var exception = new InvalidOperationException("Error handler failure.");
        await using var enumerator = reader.ReadLinesAsync(_ => throw exception, CancellationToken.None).GetAsyncEnumerator();
 
        var thrown = await Assert.ThrowsAsync<InvalidOperationException>(async () => await enumerator.MoveNextAsync());
 
        Assert.Same(exception, thrown);
    }
 
    [Fact]
    public async Task ReadLines_CompleteResponse_YieldsLines()
    {
        using var stream = new MemoryStream(Encoding.UTF8.GetBytes("first\n\nsecond\n"));
        using var reader = new StreamReader(stream);
        var lines = new List<string>();
 
        await foreach (var line in reader.ReadLinesAsync(CancellationToken.None))
        {
            lines.Add(line);
        }
 
        Assert.Equal(["first", "second"], lines);
    }
 
    private async Task<(int ExitCode, TestInteractionService InteractionService)> InvokeAsync(
        string commandName, bool standalone, bool follow, Exception exception, string data = "{}\n")
    {
        using var workspace = TemporaryWorkspace.CreateForCli(outputHelper);
        var interactionService = new TestInteractionService();
        var monitor = new TestAuxiliaryBackchannelMonitor();
        monitor.AddConnection("test", new TestAppHostAuxiliaryBackchannel
        {
            AppHostInfo = new AppHostInformation
            {
                AppHostPath = Path.Combine(workspace.WorkspaceRoot.FullName, "Test.AppHost.csproj"),
                ProcessId = Environment.ProcessId
            },
            DashboardInfoResponse = new GetDashboardInfoResponse
            {
                ApiBaseUrl = "http://localhost:18888",
                ApiToken = "test-token",
                DashboardUrls = ["http://localhost:18888"],
                IsHealthy = true
            }
        });
 
        using var handler = new MockHttpMessageHandler(request =>
        {
            if (request.RequestUri!.AbsolutePath == "/api/telemetry/resources")
            {
                return new HttpResponseMessage(HttpStatusCode.OK)
                {
                    Content = new StringContent("[]", Encoding.UTF8, "application/json")
                };
            }
 
            var content = new StreamContent(new FaultingReadStream(Encoding.UTF8.GetBytes(data), exception));
            content.Headers.ContentType = new MediaTypeHeaderValue("application/x-ndjson");
            return new HttpResponseMessage(HttpStatusCode.OK) { Content = content };
        });
 
        var services = CliTestHelper.CreateServiceCollection(workspace, outputHelper, options =>
        {
            options.InteractionServiceFactory = _ => interactionService;
            options.AuxiliaryBackchannelMonitorFactory = _ => monitor;
        });
        services.Replace(ServiceDescriptor.Singleton<IHttpClientFactory>(new MockHttpClientFactory(handler)));
        using var provider = services.BuildServiceProvider();
 
        var command = provider.GetRequiredService<RootCommand>();
        var parsed = command.Parse($"otel {commandName} --format Json{(follow ? " --follow" : "")}{(standalone ? " --dashboard-url http://localhost:18888" : "")}");
        var exitCode = await parsed.InvokeAsync().DefaultTimeout();
 
        return (exitCode, interactionService);
    }
}