File: TerminalHostAppTests.cs
Web Access
Project: src\tests\Aspire.TerminalHost.Tests\Aspire.TerminalHost.Tests.csproj (Aspire.TerminalHost.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.Sockets;
using System.Security.AccessControl;
using System.Security.Principal;
using System.Text;
using System.Text.Json;
using Aspire.Shared;
using Aspire.Shared.TerminalHost;
using Hex1b;
using Hex1b.Automation;
using Hex1b.Reflow;
using Microsoft.Extensions.Logging.Abstractions;
using StreamJsonRpc;
 
namespace Aspire.TerminalHost.Tests;
 
[CollectionDefinition(nameof(TerminalHostAppTestsCollection), DisableParallelization = true)]
public sealed class TerminalHostAppTestsCollection;
 
[Collection(nameof(TerminalHostAppTestsCollection))]
public class TerminalHostAppTests(ITestOutputHelper outputHelper)
{
    private TemporaryWorkspace CreateSocketWorkspace()
    {
        // The default workspace includes the assembly name and "Workspace" segments,
        // which can exceed macOS's 104-byte UDS path limit before the socket name is added.
        return new TemporaryWorkspace(outputHelper, Directory.CreateTempSubdirectory());
    }
 
    /// <summary>
    /// Builds a single-replica argument set for the host. Each terminal host process
    /// serves exactly one replica, so the AppHost (and these tests) just hand it one
    /// producer/consumer/control UDS path triple. The replica index is opaque to the
    /// host — callers encode it however they like in the path layout.
    /// </summary>
    private (TerminalHostArgs args, TemporaryWorkspace workspace, string controlPath) BuildArgs(
        int? columns = null,
        int? rows = null)
    {
        var workspace = CreateSocketWorkspace();
        var dcpDir = Path.Combine(workspace.Path, "terminals");
        var hostDir = dcpDir;
        var ctrlDir = dcpDir;
        SocketPermissionHelper.CreateDirectory(dcpDir, repairExisting: false);
 
        var producer = Path.Combine(dcpDir, "p.sock");
        var consumer = Path.Combine(hostDir, "r.sock");
        var control = Path.Combine(ctrlDir, "c.sock");
 
        var commandLine = new List<string>
        {
            "--producer-uds", producer,
            "--consumer-uds", consumer,
            "--control-uds", control,
        };
        if (columns is not null)
        {
            commandLine.AddRange(["--columns", columns.Value.ToString(System.Globalization.CultureInfo.InvariantCulture)]);
        }
        if (rows is not null)
        {
            commandLine.AddRange(["--rows", rows.Value.ToString(System.Globalization.CultureInfo.InvariantCulture)]);
        }
 
        var args = TerminalHostArgs.Parse([.. commandLine]);
 
        return (args, workspace, control);
    }
 
    [Fact]
    public async Task RunAsyncBindsControlListenerWhenStarted()
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
            Assert.True(File.Exists(control));
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task ControlEndpointReturnsSessionInfo()
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            using var rpc = await OpenControlRpcAsync(control);
            var info = await rpc.InvokeAsync<TerminalHostInfoResponse>(
                TerminalHostControlProtocol.GetInfoMethod);
            var session = await rpc.InvokeAsync<TerminalHostSessionInfo>(
                TerminalHostControlProtocol.GetSessionMethod);
 
            Assert.Equal(TerminalHostControlProtocol.ProtocolVersion, info.ProtocolVersion);
            Assert.Equal(args.ProducerUdsPath, session.ProducerUdsPath);
            Assert.Equal(args.ConsumerUdsPath, session.ConsumerUdsPath);
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task ShutdownRequestCausesRunAsyncToReturn()
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
        using (var rpc = await OpenControlRpcAsync(control))
        {
            // Fire and forget — the host may close the socket before the RPC ack arrives.
            _ = rpc.InvokeAsync(TerminalHostControlProtocol.ShutdownMethod);
        }
 
        var exitCode = await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        Assert.Equal(0, exitCode);
    }
 
    [Fact]
    public async Task ConcurrentControlConnectsAreRefusedDownToOne()
    {
        // The control protocol is documented as a single AppHost client (lifecycle,
        // shutdown, stats). Regression for the accept-loop race where two fast
        // back-to-back connects could both observe an empty slot before either
        // ServeClientAsync had registered into _activeRpcs, ending up with two
        // concurrently-served sessions instead of one. The reservation counter
        // increment must happen synchronously under _gate at the accept site.
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            // Open many concurrent sockets to maximise the chance of two reaching
            // the accept loop's reservation check at the same time. Whatever the
            // scheduling, exactly one must end up usable for an RPC; the rest
            // must fail (server-closed before the request returns or no response).
            const int concurrency = 8;
            var rawSockets = new Socket[concurrency];
            var connectTasks = new Task[concurrency];
            for (var i = 0; i < concurrency; i++)
            {
                rawSockets[i] = new Socket(AddressFamily.Unix, SocketType.Stream, ProtocolType.Unspecified);
                connectTasks[i] = rawSockets[i].ConnectAsync(new UnixDomainSocketEndPoint(control));
            }
 
            try
            {
                await Task.WhenAll(connectTasks).WaitAsync(TimeSpan.FromSeconds(10));
 
                // Drive an RPC over each socket and count the survivors. Refused
                // sockets either close immediately (read returns 0) or fail the
                // header-delimited read; the StreamJsonRpc completion task surfaces
                // either as a ConnectionLostException / IOException.
                var results = await Task.WhenAll(rawSockets.Select(TryGetInfoAsync));
 
                var successes = results.Count(static r => r);
                Assert.Equal(1, successes);
            }
            finally
            {
                foreach (var s in rawSockets)
                {
                    try { s.Dispose(); } catch { }
                }
            }
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
 
        static async Task<bool> TryGetInfoAsync(Socket socket)
        {
            // Each candidate gets a short window to either ack a GetInfo or be
            // refused. The refused side observes a zero-length read on its
            // NetworkStream which StreamJsonRpc raises as a ConnectionLostException.
            try
            {
                var stream = new NetworkStream(socket, ownsSocket: false);
                using var rpc = new JsonRpc(new HeaderDelimitedMessageHandler(stream, stream, new SystemTextJsonFormatter()));
                rpc.StartListening();
 
                using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(3));
                _ = await rpc.InvokeAsync<TerminalHostInfoResponse>(TerminalHostControlProtocol.GetInfoMethod).WaitAsync(cts.Token);
                return true;
            }
            catch
            {
                return false;
            }
        }
    }
 
    [Fact]
    public async Task SnapshotSessionReportsConfiguredPaths()
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            var snap = app.SnapshotSession();
            Assert.Equal(args.ProducerUdsPath, snap.ProducerUdsPath);
            Assert.Equal(args.ConsumerUdsPath, snap.ConsumerUdsPath);
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task RunAsyncWithBadArgsViaStaticEntryPointReturnsExUsage()
    {
        var exitCode = await TerminalHostApp.RunAsync(["--bogus"], CancellationToken.None);
        Assert.Equal(64, exitCode); // EX_USAGE
    }
 
    [Fact]
    public async Task HostStartsCleanlyWithStaleProducerAndConsumerSocketFiles()
    {
        // Regression: when a previous host crashes (or DCP kills the process before it
        // can gracefully unlink its UDS), the next host run hits "address already in use"
        // on Bind unless we explicitly pre-delete the path. Hex1b doesn't do this for us
        // because it doesn't know the path was previously bound by an Aspire host vs
        // some other process. Symmetry with TerminalHostControlListener (which has
        // always pre-deleted its control path).
        if (OperatingSystem.IsWindows())
        {
            // UDS files on Windows behave differently — File.Create at the path doesn't
            // produce a regular file that EADDRINUSEs the next bind, so the scenario
            // doesn't reproduce. The pre-delete still runs as defensive code.
            return;
        }
 
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        // Pre-place leftovers exactly as a crashed previous host would have left them.
        File.WriteAllBytes(args.ProducerUdsPath, []);
        File.WriteAllBytes(args.ConsumerUdsPath, []);
        Assert.True(File.Exists(args.ProducerUdsPath));
        Assert.True(File.Exists(args.ConsumerUdsPath));
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            // If pre-delete is missing, Bind fails inside the recycle loop and the
            // control listener may still come up (it pre-deletes its own path), but a
            // real producer connect would never succeed. Use the producer dial as the
            // end-to-end signal that both UDS paths are usable.
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            await using var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(10));
            await producer.SendHelloAsync(80, 24, default);
            await WaitForAsync(
                () => app.SnapshotSession().ProducerConnected,
                TimeSpan.FromSeconds(5),
                "Producer should connect after stale UDS files are pre-cleaned.");
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task SessionRecyclesAfterProducerDisconnect()
    {
        // End-to-end check of the recycle loop: the host should stay running
        // across a producer disconnect and accept a fresh producer on the
        // same UDS path, with ProducerConnected and RestartCount tracking
        // each cycle. This exercises the path DCP exercises in production
        // when the underlying process exits and gets relaunched.
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            // Initial state: host is up but no producer has dialed in yet.
            var initial = app.SnapshotSession();
            Assert.False(initial.ProducerConnected, "Session should report no producer before any connect.");
            Assert.False(initial.IsAlive, "Legacy IsAlive should mirror ProducerConnected.");
            Assert.Equal(0, initial.RestartCount);
 
            // Cycle 1: connect, get accepted, disconnect.
            await using (var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(5)))
            {
                await producer.SendHelloAsync(80, 24, default);
                await producer.SendOutputAsync("first cycle"u8.ToArray(), default);
                await WaitForAsync(
                    () => app.SnapshotSession().ProducerConnected,
                    TimeSpan.FromSeconds(5),
                    "ProducerConnected should flip to true after producer dials in.");
            }
 
            await WaitForAsync(
                () =>
                {
                    var s = app.SnapshotSession();
                    return !s.ProducerConnected && s.RestartCount >= 1;
                },
                TimeSpan.FromSeconds(10),
                "After producer disconnect, ProducerConnected should clear and RestartCount should advance.");
 
            var afterCycle1 = app.SnapshotSession();
            Assert.Equal(1, afterCycle1.RestartCount);
 
            // Cycle 2: a fresh producer should be able to dial the same UDS path.
            // This is the critical DCP-restart scenario.
            await using (var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(10)))
            {
                await producer.SendHelloAsync(80, 24, default);
                await producer.SendOutputAsync("second cycle"u8.ToArray(), default);
                await WaitForAsync(
                    () => app.SnapshotSession().ProducerConnected,
                    TimeSpan.FromSeconds(5),
                    "ProducerConnected should flip true again after the second producer dials in.");
            }
 
            await WaitForAsync(
                () =>
                {
                    var s = app.SnapshotSession();
                    return !s.ProducerConnected && s.RestartCount >= 2;
                },
                TimeSpan.FromSeconds(10),
                "After the second producer disconnects, ProducerConnected should clear and RestartCount should reach 2.");
 
            // Session itself is still there — IsAlive/ProducerConnected being false
            // is transient, the snapshot continues to report the same UDS paths.
            var afterCycle2 = app.SnapshotSession();
            Assert.Equal(args.ProducerUdsPath, afterCycle2.ProducerUdsPath);
            Assert.Equal(args.ConsumerUdsPath, afterCycle2.ConsumerUdsPath);
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task GracefulCancellationDeletesProducerAndConsumerSockets()
    {
        // Regression for https://github.com/microsoft/aspire/issues/19302: on a graceful
        // `aspire stop`, DCP signals the terminal host (SIGTERM) which Program.cs turns into a
        // cancellation of the token passed to RunAsync. That cancellation MUST flow through
        // TearDownAsync -> TerminalReplica.DisposeAsync and unlink both listen sockets
        // ({id}.dcp.sock producer + {id}.host.sock consumer). If it doesn't (the old SIGINT-only
        // behavior), the child is SIGKILLed while those sockets are still bound and both files
        // leak on disk. This test drives exactly the token-cancel path the SIGTERM handler now
        // invokes and asserts both files are gone afterward.
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            // Both listen sockets are bound while the replica waits for its first producer — this
            // is precisely the on-disk state that leaks under SIGKILL. The producer UDS is a
            // single-accept listener (its file disappears once a producer dials in), so we assert
            // on the pre-connect bound state rather than connecting a producer. Wait for both to
            // exist first; otherwise the "deleted after shutdown" assertion below could pass
            // trivially because the sockets were never bound.
            await WaitForFileAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(10));
            await WaitForFileAsync(args.ConsumerUdsPath, TimeSpan.FromSeconds(10));
            Assert.True(File.Exists(args.ProducerUdsPath), "Producer socket should be bound while the replica waits.");
            Assert.True(File.Exists(args.ConsumerUdsPath), "Consumer socket should be bound while the replica waits.");
        }
        finally
        {
            // Cancel ONLY the external token — this is precisely what Program.cs's SIGINT/SIGTERM
            // handler does (cts.Cancel()). We deliberately do NOT call app.RequestShutdown() so the
            // test exercises the signal-driven graceful path end to end.
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
 
        Assert.False(
            File.Exists(args.ProducerUdsPath),
            $"Producer socket '{args.ProducerUdsPath}' should be deleted after graceful shutdown.");
        Assert.False(
            File.Exists(args.ConsumerUdsPath),
            $"Consumer socket '{args.ConsumerUdsPath}' should be deleted after graceful shutdown.");
    }
 
    [Fact]
    public async Task SessionSnapshotIncludesNewFields()
    {
        // Even before any producer has connected, the snapshot must populate
        // the new fields so older AppHost wire deserialisation never sees a
        // missing-required-property error.
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            var snap = app.SnapshotSession();
            Assert.False(snap.ProducerConnected);
            Assert.False(snap.IsAlive);
            Assert.Equal(0, snap.RestartCount);
            Assert.Null(snap.ExitCode);
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task ConfiguredDimensionsAreAppliedUpstreamAndReportedToConsumers()
    {
        const int configuredWidth = 137;
        const int configuredHeight = 41;
        var (args, workspace, control) = BuildArgs(configuredWidth, configuredHeight);
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            await using var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(5));
            await producer.SendHelloAsync(80, 24, default);
 
            await WaitForAsync(
                () => app.SnapshotSession().ProducerConnected,
                TimeSpan.FromSeconds(5),
                "ProducerConnected should flip to true after producer dials in.");
 
            const byte FrameResize = 0x05;
            using var frameCts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
            var (type, payload) = await producer.ReadFrameAsync(frameCts.Token);
 
            Assert.Equal(FrameResize, type);
            Assert.Equal(configuredWidth, BitConverter.ToInt32(payload, 0));
            Assert.Equal(configuredHeight, BitConverter.ToInt32(payload, 4));
 
            await WaitForFileAsync(args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await using var consumer = await TestHmp1Consumer.ConnectAsync(
                args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await consumer.SendClientHelloAsync("test-consumer", "secondary", default);
            var helloPayload = await consumer.ReceiveHandshakeAsync(TimeSpan.FromSeconds(5));
 
            using var hello = JsonDocument.Parse(helloPayload);
            Assert.Equal(configuredWidth, hello.RootElement.GetProperty("width").GetInt32());
            Assert.Equal(configuredHeight, hello.RootElement.GetProperty("height").GetInt32());
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task DownstreamResizeReflowsOutputAndRetainsProducerHistory()
    {
        var (args, workspace, control) = BuildArgs(80, 24);
        using var disp = workspace;
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(30));
        var hostTask = app.RunAsync(hostCts.Token);
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
            await using var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(5));
            await producer.SendHelloAsync(80, 24, timeout.Token);
            await WaitForFileAsync(args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await using var consumer = new Hmp1WorkloadAdapter(new Hmp1ClientOptions
            {
                StreamFactory = ct => Hmp1Transports.ConnectUnixSocket(args.ConsumerUdsPath, ct),
                DefaultRole = Hmp1Role.Secondary
            });
            await consumer.ConnectAsync(timeout.Token);
            await using var mirror = Hex1bTerminal.CreateBuilder()
                .WithHeadless()
                .WithWorkload(consumer)
                .WithReflow(GhosttyReflowStrategy.Instance)
                .WithScrollback(10000)
                .Build();
            var lines = Enumerable.Range(0, 7).Select(i => $"{i}:" + new string('x', 63) + "-END").ToArray();
            await producer.SendOutputAsync(Encoding.UTF8.GetBytes(string.Join("\r\n", lines) + "\r\nready"), timeout.Token);
            await new Hex1bTerminalAutomator(mirror, TimeSpan.FromSeconds(10)).WaitUntilTextAsync("ready").WaitAsync(timeout.Token);
 
            foreach (var (width, height) in new[] { (20, 4), (80, 24) })
            {
                await consumer.RequestPrimaryAsync(width, height, timeout.Token);
                using var resized = await new Hex1bTerminalInputSequenceBuilder()
                    .WaitUntil(snapshot => snapshot.Width == width && snapshot.Height == height,
                        TimeSpan.FromSeconds(10), "The terminal host did not resize.")
                    .Build().ApplyAsync(mirror, timeout.Token);
            }
 
            var expected = string.Join('\n', lines.Select(line => line.PadRight(80)).Append("ready"));
            using var restored = await new Hex1bTerminalInputSequenceBuilder()
                .WaitUntil(snapshot => snapshot.GetScreenText().TrimEnd() == expected,
                    TimeSpan.FromSeconds(10), "The terminal host did not restore reflowed history.")
                .Build().ApplyAsync(mirror, timeout.Token);
            Assert.Equal(expected, restored.GetScreenText().TrimEnd());
 
            // A fresh peer proves the producer retained the content, not just the existing mirror.
            await using var lateConsumer = await TestHmp1Consumer.ConnectAsync(args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await lateConsumer.SendClientHelloAsync("late-reflow-peer", "secondary", timeout.Token);
            await lateConsumer.ReceiveHandshakeAsync(TimeSpan.FromSeconds(5));
            var replayWorkload = new Hex1bAppWorkloadAdapter();
            await using var replay = Hex1bTerminal.CreateBuilder()
                .WithHeadless().WithDimensions(80, 24).WithWorkload(replayWorkload).Build();
            replayWorkload.Write(Encoding.UTF8.GetString(lateConsumer.InitialState) + "\r\nreplay-complete");
            using var lateSnapshot = await new Hex1bTerminalInputSequenceBuilder()
                .WaitUntil(snapshot => snapshot.ContainsText("replay-complete"),
                    TimeSpan.FromSeconds(10), "The late peer did not consume its initial state.")
                .Build().ApplyAsync(replay, timeout.Token);
            Assert.Equal(expected + new string(' ', 80 - "ready".Length) + "\nreplay-complete",
                lateSnapshot.GetScreenText().TrimEnd());
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task DownstreamPrimaryResizeIsForwardedUpstreamAsRawResizeFrame()
    {
        // Regression: the consumer-side multi-head server fires its OnResized
        // event whenever the current primary peer changes the producer dims
        // (RequestPrimary or explicit Resize from primary). The terminal host
        // bridges that event to a raw HMP1 FrameResize (0x05) on the upstream
        // (DCP-facing) connection — bypassing Hex1b's stock Hmp1WorkloadAdapter
        // IsPrimary gate, which would silently drop every resize because DCP's
        // minimal HMP1 server never sends Hello.PrimaryPeerId or RoleChange.
        // Without this bridge, the underlying PTY stayed at its DCP-initial
        // dims forever.
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            // Stand up the upstream "DCP" first. Send Hello so the host's
            // upstream-side adapter handshakes successfully and ProducerConnected
            // flips before we attempt the consumer-side connect — otherwise the
            // resize broadcast would race the upstream stream's existence.
            await using var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(5));
            await producer.SendHelloAsync(80, 24, default);
 
            await WaitForAsync(
                () => app.SnapshotSession().ProducerConnected,
                TimeSpan.FromSeconds(5),
                "ProducerConnected should flip to true after producer dials in.");
 
            // Wait for the consumer-side UDS server to bind so the dial below
            // doesn't race a not-yet-listening socket.
            await WaitForFileAsync(args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
 
            // Now connect a minimal raw-frame HMP1 client to the consumer UDS:
            // ClientHello + RequestPrimary, then keep the stream open. We don't
            // need a full Hex1bTerminal because all we're verifying is that the
            // server-side OnResized event (which fires when RequestPrimary
            // promotes us and applies the requested dims) is bridged upstream.
            const int requestedWidth = 123;
            const int requestedHeight = 45;
 
            await using var consumer = await TestHmp1Consumer.ConnectAsync(
                args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await consumer.SendClientHelloAsync("test-consumer", "primary", default);
            await consumer.ReceiveHandshakeAsync(TimeSpan.FromSeconds(5));
            await consumer.SendRequestPrimaryAsync(requestedWidth, requestedHeight, default);
 
            // The frame the test producer should observe upstream:
            // [type=0x05 Resize][length=8 LE][width:4B LE][height:4B LE]
            //
            // Use the predicate variant because the host's Hex1bTerminal may
            // emit additional resize frames upstream during startup or in
            // response to incoming Hello dims (the consumer-side server's
            // Resized event also triggers Hex1bTerminal's own
            // workload.ResizeAsync). We only care that SOME FrameResize lands
            // upstream with the dims the consumer requested via RequestPrimary.
            const byte FrameResize = 0x05;
            var payload = await producer.WaitForMatchingFrameAsync(
                FrameResize,
                p =>
                {
                    if (p.Length != 8)
                    {
                        return false;
                    }
                    var w = p[0] | (p[1] << 8) | (p[2] << 16) | (p[3] << 24);
                    var h = p[4] | (p[5] << 8) | (p[6] << 16) | (p[7] << 24);
                    return w == requestedWidth && h == requestedHeight;
                },
                TimeSpan.FromSeconds(10));
 
            // Sanity: payload encodes exactly the requested dims (defensive
            // — the predicate already filtered, but keeps the assertion intent
            // explicit on the test surface).
            Assert.Equal(8, payload.Length);
            var observedWidth = payload[0] | (payload[1] << 8) | (payload[2] << 16) | (payload[3] << 24);
            var observedHeight = payload[4] | (payload[5] << 8) | (payload[6] << 16) | (payload[7] << 24);
            Assert.Equal(requestedWidth, observedWidth);
            Assert.Equal(requestedHeight, observedHeight);
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task ConsumerListenerBindFailureIsReportedBeforeTerminalStarts()
    {
        using var workspace = CreateSocketWorkspace();
        Directory.CreateDirectory(Path.Combine(workspace.Path, ".aspire"));
        var parentFile = Path.Combine(workspace.Path, ".aspire", "trmnl");
        await File.WriteAllTextAsync(parentFile, "");
        await using var presentation = new Hmp1PresentationAdapter();
 
        Assert.Throws<IOException>(() => new Hmp1UdsServerListenerFilter(
            Path.Combine(parentFile, "consumer.sock"),
            presentation,
            NullLogger<Hmp1UdsServerListenerFilter>.Instance,
            _ => { }));
    }
 
    [Fact]
    public async Task ConsumerListenerWaitsForAcceptedClientHandshakeDuringTeardown()
    {
        using var workspace = CreateSocketWorkspace();
        var socketPath = Path.Combine(workspace.Path, ".aspire", "trmnl", "consumer.sock");
        await using var presentation = new Hmp1PresentationAdapter();
        var callbackStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
        var releaseCallback = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
        presentation.OnClientConnected = async (_, _) =>
        {
            callbackStarted.TrySetResult();
            await releaseCallback.Task.ConfigureAwait(false);
        };
 
        using var listener = new Hmp1UdsServerListenerFilter(
            socketPath,
            presentation,
            NullLogger<Hmp1UdsServerListenerFilter>.Instance,
            _ => { });
        await listener.OnSessionStartAsync(80, 24, DateTimeOffset.UtcNow);
 
        Task? endTask = null;
        try
        {
            await using var consumer = await TestHmp1Consumer.ConnectAsync(socketPath, TimeSpan.FromSeconds(5));
            await consumer.SendClientHelloAsync("test-consumer", "secondary", default);
            await callbackStarted.Task.WaitAsync(TimeSpan.FromSeconds(5));
 
            endTask = listener.OnSessionEndAsync(TimeSpan.Zero).AsTask();
            var completed = await Task.WhenAny(endTask, Task.Delay(TimeSpan.FromMilliseconds(200)));
            Assert.NotSame(endTask, completed);
        }
        finally
        {
            releaseCallback.TrySetResult();
            if (endTask is null)
            {
                await listener.OnSessionEndAsync(TimeSpan.Zero);
            }
            else
            {
                await endTask.WaitAsync(TimeSpan.FromSeconds(5));
            }
        }
    }
 
    [Fact]
    public async Task CurrentDimensionsArePreservedAcrossProducerRecycle()
    {
        const int requestedWidth = 123;
        const int requestedHeight = 45;
        const byte frameResize = 0x05;
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        using var hostCts = new CancellationTokenSource();
        var hostTask = app.RunAsync(hostCts.Token);
 
        try
        {
            await WaitForFileAsync(control, TimeSpan.FromSeconds(10));
 
            await using (var producer = await ConnectProducerAsync(args.ProducerUdsPath, TimeSpan.FromSeconds(5)))
            {
                await producer.SendHelloAsync(80, 24, default);
                await WaitForAsync(
                    () => app.SnapshotSession().ProducerConnected,
                    TimeSpan.FromSeconds(5),
                    "The first producer should connect.");
 
                await using var consumer = await TestHmp1Consumer.ConnectAsync(
                    args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
                await consumer.SendClientHelloAsync("test-consumer", "primary", default);
                await consumer.ReceiveHandshakeAsync(TimeSpan.FromSeconds(5));
                await consumer.SendRequestPrimaryAsync(requestedWidth, requestedHeight, default);
 
                _ = await producer.WaitForMatchingFrameAsync(
                    frameResize,
                    payload => payload.Length == 8
                        && BitConverter.ToInt32(payload, 0) == requestedWidth
                        && BitConverter.ToInt32(payload, 4) == requestedHeight,
                    TimeSpan.FromSeconds(10));
                await WaitForAsync(
                    () =>
                    {
                        var session = app.SnapshotSession();
                        return session.CurrentColumns == requestedWidth && session.CurrentRows == requestedHeight;
                    },
                    TimeSpan.FromSeconds(5),
                    "The resized dimensions should become authoritative.");
            }
 
            await WaitForAsync(
                () => app.SnapshotSession().RestartCount >= 1,
                TimeSpan.FromSeconds(10),
                "The terminal should recycle after the first producer disconnects.");
 
            await using var replacementProducer = await ConnectProducerAsync(
                args.ProducerUdsPath, TimeSpan.FromSeconds(10));
            await replacementProducer.SendHelloAsync(80, 24, default);
 
            _ = await replacementProducer.WaitForMatchingFrameAsync(
                frameResize,
                payload => payload.Length == 8
                    && BitConverter.ToInt32(payload, 0) == requestedWidth
                    && BitConverter.ToInt32(payload, 4) == requestedHeight,
                TimeSpan.FromSeconds(10));
 
            await using var replacementConsumer = await TestHmp1Consumer.ConnectAsync(
                args.ConsumerUdsPath, TimeSpan.FromSeconds(5));
            await replacementConsumer.SendClientHelloAsync("replacement-consumer", "secondary", default);
            var helloPayload = await replacementConsumer.ReceiveHandshakeAsync(TimeSpan.FromSeconds(5));
 
            using var hello = JsonDocument.Parse(helloPayload);
            Assert.Equal(requestedWidth, hello.RootElement.GetProperty("width").GetInt32());
            Assert.Equal(requestedHeight, hello.RootElement.GetProperty("height").GetInt32());
        }
        finally
        {
            app.RequestShutdown();
            hostCts.Cancel();
            await hostTask.WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    /// <summary>
    /// A minimal HMP1 server-role producer for tests. Connects to the UDS path
    /// the terminal host is listening on (producer side) and writes the bare
    /// minimum frames the Hex1bTerminal client expects: a Hello, optional
    /// Output frames, and EOF on dispose.
    /// </summary>
    private sealed class TestHmp1Producer : IAsyncDisposable
    {
        // HMP1 wire format: [type:1B][length:4B LE][payload:N bytes].
        private const byte FrameHello = 0x01;
        private const byte FrameOutput = 0x03;
 
        private readonly Socket _socket;
        private readonly NetworkStream _stream;
        private bool _disposed;
 
        public TestHmp1Producer(Socket socket)
        {
            _socket = socket;
            _stream = new NetworkStream(socket, ownsSocket: true);
        }
 
        public async Task SendHelloAsync(int width, int height, CancellationToken ct)
        {
            var json = $"{{\"version\":1,\"width\":{width},\"height\":{height}}}";
            await SendFrameAsync(FrameHello, System.Text.Encoding.UTF8.GetBytes(json), ct).ConfigureAwait(false);
        }
 
        public Task SendOutputAsync(byte[] payload, CancellationToken ct) =>
            SendFrameAsync(FrameOutput, payload, ct);
 
        /// <summary>
        /// Drains HMP1 frames until a frame of the requested type whose
        /// payload passes the <paramref name="match"/> predicate is found.
        /// Useful when upstream may receive multiple frames of the same type
        /// from different sources (e.g., terminal startup resize vs explicit
        /// user resize).
        /// </summary>
        public async Task<byte[]> WaitForMatchingFrameAsync(
            byte expectedType, Func<byte[], bool> match, TimeSpan timeout)
        {
            using var cts = new CancellationTokenSource(timeout);
            while (true)
            {
                var (type, payload) = await ReadFrameAsync(cts.Token).ConfigureAwait(false);
                if (type == expectedType && match(payload))
                {
                    return payload;
                }
            }
        }
 
        /// <summary>Reads exactly one HMP1 frame and returns its (type, payload).</summary>
        public async Task<(byte Type, byte[] Payload)> ReadFrameAsync(CancellationToken ct)
        {
            var header = new byte[5];
            await ReadExactlyAsync(header, ct).ConfigureAwait(false);
 
            var type = header[0];
            var length = (uint)(header[1] | (header[2] << 8) | (header[3] << 16) | (header[4] << 24));
            if (length > 16 * 1024 * 1024)
            {
                throw new InvalidOperationException($"Producer-side reader: frame length {length} exceeds 16MB cap.");
            }
 
            var payload = new byte[length];
            if (length > 0)
            {
                await ReadExactlyAsync(payload, ct).ConfigureAwait(false);
            }
            return (type, payload);
        }
 
        private async Task ReadExactlyAsync(byte[] buffer, CancellationToken ct)
        {
            var offset = 0;
            while (offset < buffer.Length)
            {
                var read = await _stream.ReadAsync(buffer.AsMemory(offset, buffer.Length - offset), ct).ConfigureAwait(false);
                if (read == 0)
                {
                    throw new EndOfStreamException(
                        $"Producer-side reader: stream EOF after {offset} of {buffer.Length} bytes.");
                }
                offset += read;
            }
        }
 
        private async Task SendFrameAsync(byte type, byte[] payload, CancellationToken ct)
        {
            var header = new byte[5];
            header[0] = type;
            header[1] = (byte)(payload.Length & 0xFF);
            header[2] = (byte)((payload.Length >> 8) & 0xFF);
            header[3] = (byte)((payload.Length >> 16) & 0xFF);
            header[4] = (byte)((payload.Length >> 24) & 0xFF);
            await _stream.WriteAsync(header, ct).ConfigureAwait(false);
            if (payload.Length > 0)
            {
                await _stream.WriteAsync(payload, ct).ConfigureAwait(false);
            }
            await _stream.FlushAsync(ct).ConfigureAwait(false);
        }
 
        public async ValueTask DisposeAsync()
        {
            if (_disposed)
            {
                return;
            }
            _disposed = true;
            try { _socket.Shutdown(SocketShutdown.Both); } catch { /* ignore */ }
            await _stream.DisposeAsync().ConfigureAwait(false);
        }
    }
 
    private static async Task<TestHmp1Producer> ConnectProducerAsync(string socketPath, TimeSpan timeout)
    {
        // Retry loop because there is a brief unbound window between recycle
        // iterations on the host side; the test producer should ride through
        // that the same way DCP does in production.
        var sw = System.Diagnostics.Stopwatch.StartNew();
        Exception? last = null;
        while (sw.Elapsed < timeout)
        {
            var socket = new Socket(AddressFamily.Unix, SocketType.Stream, ProtocolType.Unspecified);
            try
            {
                await socket.ConnectAsync(new UnixDomainSocketEndPoint(socketPath)).ConfigureAwait(false);
                return new TestHmp1Producer(socket);
            }
            catch (Exception ex)
            {
                socket.Dispose();
                last = ex;
                await Task.Delay(50).ConfigureAwait(false);
            }
        }
        throw new TimeoutException(
            $"Timed out connecting to producer UDS '{socketPath}' after {timeout.TotalSeconds:F1}s.", last);
    }
 
    /// <summary>
    /// A minimal HMP1 client-role consumer for tests. Connects to the consumer
    /// UDS the terminal host is listening on, then writes raw HMP1 frames
    /// (ClientHello, RequestPrimary). Avoids spinning up a full
    /// <c>Hex1bTerminal</c> in a test process where no interactive console is
    /// attached.
    /// </summary>
    private sealed class TestHmp1Consumer : IAsyncDisposable
    {
        private const byte FrameHello = 0x01;
        private const byte FrameStateSync = 0x02;
        private const byte FrameRequestPrimary = 0x07;
        private const byte FrameClientHello = 0x0B;
 
        private readonly Socket _socket;
        private readonly NetworkStream _stream;
        private bool _disposed;
 
        public byte[] InitialState { get; private set; } = [];
 
        private TestHmp1Consumer(Socket socket)
        {
            _socket = socket;
            _stream = new NetworkStream(socket, ownsSocket: true);
        }
 
        public static async Task<TestHmp1Consumer> ConnectAsync(string socketPath, TimeSpan timeout)
        {
            var sw = System.Diagnostics.Stopwatch.StartNew();
            Exception? last = null;
            while (sw.Elapsed < timeout)
            {
                var socket = new Socket(AddressFamily.Unix, SocketType.Stream, ProtocolType.Unspecified);
                try
                {
                    await socket.ConnectAsync(new UnixDomainSocketEndPoint(socketPath)).ConfigureAwait(false);
                    return new TestHmp1Consumer(socket);
                }
                catch (Exception ex)
                {
                    socket.Dispose();
                    last = ex;
                    await Task.Delay(50).ConfigureAwait(false);
                }
            }
            throw new TimeoutException(
                $"Timed out connecting to consumer UDS '{socketPath}' after {timeout.TotalSeconds:F1}s.", last);
        }
 
        public async Task SendClientHelloAsync(string displayName, string defaultRole, CancellationToken ct)
        {
            // JSON keys are camelCase per Hmp1JsonContext.PropertyNamingPolicy.
            // Roles on the wire are the lowercase strings "primary" / "secondary".
            var json = $"{{\"displayName\":\"{displayName}\",\"defaultRole\":\"{defaultRole}\"}}";
            await SendFrameAsync(FrameClientHello, System.Text.Encoding.UTF8.GetBytes(json), ct).ConfigureAwait(false);
        }
 
        public async Task<byte[]> ReceiveHandshakeAsync(TimeSpan timeout)
        {
            using var cts = new CancellationTokenSource(timeout);
            var (helloType, helloPayload) = await ReadFrameAsync(cts.Token).ConfigureAwait(false);
            var (stateSyncType, stateSyncPayload) = await ReadFrameAsync(cts.Token).ConfigureAwait(false);
 
            if (helloType != FrameHello || stateSyncType != FrameStateSync)
            {
                throw new InvalidDataException(
                    $"Expected Hello and StateSync frames, received 0x{helloType:X2} and 0x{stateSyncType:X2}.");
            }
 
            InitialState = stateSyncPayload;
            return helloPayload;
        }
 
        public async Task SendRequestPrimaryAsync(int cols, int rows, CancellationToken ct)
        {
            var json = $"{{\"cols\":{cols},\"rows\":{rows}}}";
            await SendFrameAsync(FrameRequestPrimary, System.Text.Encoding.UTF8.GetBytes(json), ct).ConfigureAwait(false);
        }
 
        private async Task<(byte Type, byte[] Payload)> ReadFrameAsync(CancellationToken ct)
        {
            // HMP1 frames are [type:1B][length:4B LE][payload:N bytes].
            var header = new byte[5];
            await ReadExactlyAsync(header, ct).ConfigureAwait(false);
 
            var length = header[1] | (header[2] << 8) | (header[3] << 16) | (header[4] << 24);
            if (length < 0 || length > 16 * 1024 * 1024)
            {
                throw new InvalidDataException($"Consumer-side reader received invalid frame length {length}.");
            }
 
            var payload = new byte[length];
            if (payload.Length > 0)
            {
                await ReadExactlyAsync(payload, ct).ConfigureAwait(false);
            }
 
            return (header[0], payload);
        }
 
        private async Task ReadExactlyAsync(byte[] buffer, CancellationToken ct)
        {
            var offset = 0;
            while (offset < buffer.Length)
            {
                var read = await _stream.ReadAsync(buffer.AsMemory(offset), ct).ConfigureAwait(false);
                if (read == 0)
                {
                    throw new EndOfStreamException(
                        $"Consumer-side reader: stream EOF after {offset} of {buffer.Length} bytes.");
                }
                offset += read;
            }
        }
 
        private async Task SendFrameAsync(byte type, byte[] payload, CancellationToken ct)
        {
            var header = new byte[5];
            header[0] = type;
            header[1] = (byte)(payload.Length & 0xFF);
            header[2] = (byte)((payload.Length >> 8) & 0xFF);
            header[3] = (byte)((payload.Length >> 16) & 0xFF);
            header[4] = (byte)((payload.Length >> 24) & 0xFF);
            await _stream.WriteAsync(header, ct).ConfigureAwait(false);
            if (payload.Length > 0)
            {
                await _stream.WriteAsync(payload, ct).ConfigureAwait(false);
            }
            await _stream.FlushAsync(ct).ConfigureAwait(false);
        }
 
        public async ValueTask DisposeAsync()
        {
            if (_disposed)
            {
                return;
            }
            _disposed = true;
            try { _socket.Shutdown(SocketShutdown.Both); } catch { /* ignore */ }
            await _stream.DisposeAsync().ConfigureAwait(false);
        }
    }
 
    private static async Task WaitForAsync(Func<bool> predicate, TimeSpan timeout, string failureMessage)
    {
        var sw = System.Diagnostics.Stopwatch.StartNew();
        while (sw.Elapsed < timeout)
        {
            if (predicate())
            {
                return;
            }
            await Task.Delay(25).ConfigureAwait(false);
        }
        throw new TimeoutException($"{failureMessage} (waited {timeout.TotalSeconds:F1}s).");
    }
 
    private static async Task<JsonRpc> OpenControlRpcAsync(string socketPath)
    {
        var socket = new Socket(AddressFamily.Unix, SocketType.Stream, ProtocolType.Unspecified);
        try
        {
            await socket.ConnectAsync(new UnixDomainSocketEndPoint(socketPath));
        }
        catch
        {
            socket.Dispose();
            throw;
        }
 
        var stream = new NetworkStream(socket, ownsSocket: true);
        var formatter = new SystemTextJsonFormatter();
        var handler = new HeaderDelimitedMessageHandler(stream, stream, formatter);
        var rpc = new JsonRpc(handler);
        rpc.StartListening();
        return rpc;
    }
 
    [Fact]
    public async Task ControlSocketIsRestrictedToOwningUser()
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
 
        await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
        await using var listener = new TerminalHostControlListener(
            control, new TerminalHostControlRpcTarget(app), NullLogger.Instance);
        await listener.StartAsync();
 
        AssertSocketIsRestrictedToOwningUser(control);
        using var rpc = await OpenControlRpcAsync(control);
        var info = await rpc.InvokeAsync<TerminalHostInfoResponse>(
            TerminalHostControlProtocol.GetInfoMethod).WaitAsync(TimeSpan.FromSeconds(10));
        Assert.Equal(TerminalHostControlProtocol.ProtocolVersion, info.ProtocolVersion);
    }
 
    [Fact]
    public async Task ProducerSocketIsRestrictedToOwningUserAcrossRebinds()
    {
        var (args, workspace, _) = BuildArgs();
        using var disp = workspace;
 
        for (var cycle = 0; cycle < 2; cycle++)
        {
            File.Delete(args.ProducerUdsPath);
            using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
            var acceptTask = TerminalReplica.AcceptProducerAsync(args.ProducerUdsPath, cts.Token);
 
            try
            {
                // AcceptProducerAsync binds and restricts the socket synchronously before
                // awaiting a client; do not poll permissions and hide an asynchronous chmod race.
                AssertSocketIsRestrictedToOwningUser(args.ProducerUdsPath);
                using var socket = new Socket(AddressFamily.Unix, SocketType.Stream, ProtocolType.Unspecified);
                await socket.ConnectAsync(new UnixDomainSocketEndPoint(args.ProducerUdsPath), cts.Token);
                await using var producer = new NetworkStream(socket, ownsSocket: false);
                await using var accepted = await acceptTask;
                await producer.WriteAsync(new byte[] { 42 }, cts.Token);
                var received = new byte[1];
                await accepted.ReadExactlyAsync(received, cts.Token);
                Assert.Equal(42, received[0]);
            }
            finally
            {
                await cts.CancelAsync();
                try
                {
                    await acceptTask;
                }
                catch (OperationCanceledException)
                {
                }
            }
        }
    }
 
    [Fact]
    public async Task ConsumerSocketIsRestrictedToOwningUser()
    {
        var (args, workspace, _) = BuildArgs();
        using var disp = workspace;
        await using var presentation = new Hmp1PresentationAdapter();
        var connected = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
        presentation.OnClientConnected = (_, _) =>
        {
            connected.TrySetResult();
            return Task.CompletedTask;
        };
 
        using var listener = new Hmp1UdsServerListenerFilter(
            args.ConsumerUdsPath,
            presentation,
            NullLogger<Hmp1UdsServerListenerFilter>.Instance,
            ex => connected.TrySetException(ex));
        AssertSocketIsRestrictedToOwningUser(args.ConsumerUdsPath);
        await listener.OnSessionStartAsync(80, 24, DateTimeOffset.UtcNow);
 
        try
        {
            await using var consumer = await TestHmp1Consumer.ConnectAsync(
                args.ConsumerUdsPath, TimeSpan.FromSeconds(10));
            using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
            await consumer.SendClientHelloAsync("owner", "secondary", cts.Token);
            await connected.Task.WaitAsync(cts.Token);
        }
        finally
        {
            await listener.OnSessionEndAsync(TimeSpan.Zero).AsTask().WaitAsync(TimeSpan.FromSeconds(10));
        }
    }
 
    [Fact]
    public async Task ProducerListenerBindFailureIsReported()
    {
        using var workspace = CreateSocketWorkspace();
        Directory.CreateDirectory(Path.Combine(workspace.Path, ".aspire"));
        var parentFile = Path.Combine(workspace.Path, ".aspire", "trmnl");
        await File.WriteAllTextAsync(parentFile, "");
 
        await Assert.ThrowsAsync<IOException>(() => TerminalReplica.AcceptProducerAsync(
            Path.Combine(parentFile, "producer.sock"), CancellationToken.None));
    }
 
    [Fact]
    public async Task ProducerListenerCanBeReboundAfterCancellation()
    {
        var (args, workspace, _) = BuildArgs();
        using var disp = workspace;
 
        for (var cycle = 0; cycle < 2; cycle++)
        {
            File.Delete(args.ProducerUdsPath);
            using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
            var acceptTask = TerminalReplica.AcceptProducerAsync(args.ProducerUdsPath, cts.Token);
            try
            {
                AssertSocketIsRestrictedToOwningUser(args.ProducerUdsPath);
            }
            finally
            {
                await cts.CancelAsync();
            }
            await Assert.ThrowsAnyAsync<OperationCanceledException>(() => acceptTask);
        }
    }
 
    [Theory]
    [InlineData("control")]
    [InlineData("producer")]
    [InlineData("consumer")]
    public async Task ListenersRejectPermissiveDirectoryWithoutChangingPermissionsOrDeletingFiles(string listenerKind)
    {
        var (args, workspace, control) = BuildArgs();
        using var disp = workspace;
        MakeSocketDirectoryPermissive(control);
        var directory = Path.GetDirectoryName(control)!;
        var originalPermissions = OperatingSystem.IsWindows()
            ? new DirectoryInfo(directory).GetAccessControl().GetSecurityDescriptorSddlForm(AccessControlSections.All)
            : File.GetUnixFileMode(directory).ToString();
        var path = listenerKind switch
        {
            "control" => control,
            "producer" => args.ProducerUdsPath,
            _ => args.ConsumerUdsPath
        };
        await File.WriteAllTextAsync(path, "not our socket");
 
        if (listenerKind == "control")
        {
            await using var app = new TerminalHostApp(args, NullLoggerFactory.Instance);
            await using var listener = new TerminalHostControlListener(
                control, new TerminalHostControlRpcTarget(app), NullLogger.Instance);
            await Assert.ThrowsAsync<IOException>(listener.StartAsync);
        }
        else if (listenerKind == "producer")
        {
            await Assert.ThrowsAsync<IOException>(() => TerminalReplica.AcceptProducerAsync(path, CancellationToken.None));
        }
        else
        {
            await using var presentation = new Hmp1PresentationAdapter();
            Assert.Throws<IOException>(() => new Hmp1UdsServerListenerFilter(
                path, presentation, NullLogger<Hmp1UdsServerListenerFilter>.Instance, _ => { }));
        }
 
        Assert.Equal("not our socket", await File.ReadAllTextAsync(path));
        var actualPermissions = OperatingSystem.IsWindows()
            ? new DirectoryInfo(directory).GetAccessControl().GetSecurityDescriptorSddlForm(AccessControlSections.All)
            : File.GetUnixFileMode(directory).ToString();
        Assert.Equal(originalPermissions, actualPermissions);
    }
 
    private static void MakeSocketDirectoryPermissive(string socketPath)
    {
        var path = Path.GetDirectoryName(socketPath)!;
        if (OperatingSystem.IsWindows())
        {
            var directory = new DirectoryInfo(path);
            var security = directory.GetAccessControl();
            security.AddAccessRule(new FileSystemAccessRule(
                new SecurityIdentifier(WellKnownSidType.WorldSid, null),
                FileSystemRights.FullControl,
                InheritanceFlags.ContainerInherit | InheritanceFlags.ObjectInherit,
                PropagationFlags.None,
                AccessControlType.Allow));
            directory.SetAccessControl(security);
            return;
        }
 
        File.SetUnixFileMode(path,
            UnixFileMode.UserRead | UnixFileMode.UserWrite | UnixFileMode.UserExecute |
            UnixFileMode.GroupRead | UnixFileMode.GroupWrite | UnixFileMode.GroupExecute |
            UnixFileMode.OtherRead | UnixFileMode.OtherWrite | UnixFileMode.OtherExecute);
    }
 
    private static void AssertSocketIsRestrictedToOwningUser(string socketPath)
    {
        var directoryPath = Path.GetDirectoryName(socketPath)!;
        if (OperatingSystem.IsWindows())
        {
            using var identity = WindowsIdentity.GetCurrent();
            var directorySecurity = new DirectoryInfo(directoryPath).GetAccessControl();
            Assert.True(directorySecurity.AreAccessRulesProtected);
            Assert.Equal(identity.User, directorySecurity.GetOwner(typeof(SecurityIdentifier)));
            FileSystemSecurity[] descriptors =
            [
                directorySecurity,
                new FileInfo(socketPath).GetAccessControl()
            ];
            foreach (var security in descriptors)
            {
                var rule = Assert.Single(security.GetAccessRules(true, true, typeof(SecurityIdentifier))
                    .Cast<FileSystemAccessRule>());
                Assert.Equal(identity.User, rule.IdentityReference);
                Assert.Equal(AccessControlType.Allow, rule.AccessControlType);
                Assert.Equal(FileSystemRights.FullControl, rule.FileSystemRights);
            }
        }
        else
        {
            Assert.Equal(
                UnixFileMode.UserRead | UnixFileMode.UserWrite | UnixFileMode.UserExecute,
                File.GetUnixFileMode(directoryPath));
            Assert.Equal(
                UnixFileMode.UserRead | UnixFileMode.UserWrite,
                File.GetUnixFileMode(socketPath));
        }
    }
 
    private static async Task WaitForFileAsync(string path, TimeSpan timeout)
    {
        var sw = System.Diagnostics.Stopwatch.StartNew();
        while (sw.Elapsed < timeout)
        {
            if (File.Exists(path))
            {
                return;
            }
            await Task.Delay(50);
        }
 
        throw new TimeoutException($"Timed out waiting for '{path}' after {timeout.TotalSeconds:F1}s.");
    }
}