1569 lines
73 KiB
C#
1569 lines
73 KiB
C#
using System.Diagnostics;
|
||
using Google.Protobuf.WellKnownTypes;
|
||
using Microsoft.Extensions.Options;
|
||
using Microsoft.Extensions.Time.Testing;
|
||
using ZB.MOM.WW.MxGateway.Contracts.Proto;
|
||
using ZB.MOM.WW.MxGateway.Server.Configuration;
|
||
using ZB.MOM.WW.MxGateway.Server.Metrics;
|
||
using ZB.MOM.WW.MxGateway.Server.Sessions;
|
||
using ZB.MOM.WW.MxGateway.Server.Workers;
|
||
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
|
||
|
||
namespace ZB.MOM.WW.MxGateway.Tests.Gateway.Sessions;
|
||
|
||
public sealed class SessionManagerTests
|
||
{
|
||
/// <summary>Verifies that opening a session with a ready worker registers the session in ready state.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_WithWorkerReady_RegistersReadySession()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
FakeSessionWorkerClientFactory factory = new(workerClient)
|
||
{
|
||
ApplyLifecycleTransitions = true,
|
||
};
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(factory, metrics: metrics);
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
Assert.True(manager.TryGetSession(session.SessionId, out GatewaySession? registered));
|
||
Assert.Same(session, registered);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
Assert.Equal("client-1", session.ClientIdentity);
|
||
Assert.Equal(["StartingWorker", "WaitingForPipe", "Handshaking", "InitializingWorker"], factory.ObservedStates);
|
||
Assert.Equal(1, metrics.GetSnapshot().OpenSessions);
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsOpened);
|
||
}
|
||
|
||
/// <summary>
|
||
/// Verifies the pipe name stays short enough that its Unix-domain-socket path
|
||
/// (TMPDIR + "CoreFxPipe_" + name) fits the 104-byte macOS sun_path limit under the
|
||
/// default per-user TMPDIR (~49 chars), and keeps the pid + session-guid uniqueness
|
||
/// contract (NEXT-01).
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_PipeNameIsShortAndUniquePerPidAndSession()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
FakeSessionWorkerClientFactory factory = new(workerClient)
|
||
{
|
||
ApplyLifecycleTransitions = true,
|
||
};
|
||
SessionManager manager = CreateManager(factory);
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
Assert.Matches($"^mxgw-{Environment.ProcessId}-[0-9a-f]{{32}}$", session.PipeName);
|
||
Assert.EndsWith(session.SessionId["session-".Length..], session.PipeName, StringComparison.Ordinal);
|
||
|
||
// 104-byte sun_path − NUL − ~49-char default macOS TMPDIR − "CoreFxPipe_".
|
||
const int MaxPipeNameLength = 104 - 1 - 49 - 11;
|
||
|
||
// macOS pids top out at 99999, so 5 digits is the worst case the budget must survive.
|
||
// The check must substitute that worst case for the *running* pid's digit count rather
|
||
// than pad the measured length upward: Windows pids are routinely 6 digits, which would
|
||
// otherwise fail this assertion on a host whose own pipe-name limit (256) is irrelevant
|
||
// to the macOS budget being guarded here.
|
||
const int WorstCaseMacOsPidDigits = 5;
|
||
int runningPidDigits = Environment.ProcessId
|
||
.ToString(System.Globalization.CultureInfo.InvariantCulture).Length;
|
||
int worstCaseLength = session.PipeName.Length - runningPidDigits + WorstCaseMacOsPidDigits;
|
||
Assert.True(
|
||
worstCaseLength <= MaxPipeNameLength,
|
||
$"Pipe name '{session.PipeName}' would overflow the macOS socket-path budget at a 5-digit pid " +
|
||
$"({worstCaseLength} > {MaxPipeNameLength}).");
|
||
}
|
||
|
||
/// <summary>Verifies that a session opened by an authenticated caller records that caller's API key id in OwnerKeyId.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_WithOwnerKeyId_RecordsOwnerKeyIdOnSession()
|
||
{
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(new FakeWorkerClient()));
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
clientIdentity: "MyKey Display",
|
||
ownerKeyId: "key-abc123",
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal("key-abc123", session.OwnerKeyId);
|
||
}
|
||
|
||
/// <summary>Verifies that a session opened without an owner key id records null in OwnerKeyId.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_WithNullOwnerKeyId_RecordsNullOwnerKeyIdOnSession()
|
||
{
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(new FakeWorkerClient()));
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
clientIdentity: null,
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Assert.Null(session.OwnerKeyId);
|
||
}
|
||
|
||
/// <summary>Verifies that opening a session sets the initial lease expiry from the configured default lease.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_SetsInitialDefaultLease()
|
||
{
|
||
ManualTimeProvider clock = new(DateTimeOffset.Parse("2026-04-29T10:00:00Z", System.Globalization.CultureInfo.InvariantCulture));
|
||
GatewayOptions options = CreateOptions(defaultLeaseSeconds: 1800);
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(new FakeWorkerClient()),
|
||
options: options,
|
||
timeProvider: clock);
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
Assert.Equal(clock.GetUtcNow() + TimeSpan.FromMinutes(30), session.LeaseExpiresAt);
|
||
}
|
||
|
||
/// <summary>Verifies that session generation creates client correlation ID from client name and session ID.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_GeneratesClientCorrelationIdFromClientNameAndSessionId()
|
||
{
|
||
SessionOpenRequest request = CreateOpenRequest() with
|
||
{
|
||
ClientSessionName = "rust-load-client",
|
||
ClientCorrelationId = "caller-provided-correlation",
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(new FakeWorkerClient()));
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(request, "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
Assert.Equal($"rust-load-client-{session.SessionId}", session.ClientCorrelationId);
|
||
}
|
||
|
||
/// <summary>Verifies that opening a session without a client session name uses the client correlation prefix.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_WhenClientSessionNameMissing_UsesClientCorrelationPrefix()
|
||
{
|
||
SessionOpenRequest request = CreateOpenRequest() with
|
||
{
|
||
ClientSessionName = "",
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(new FakeWorkerClient()));
|
||
|
||
GatewaySession session = await manager.OpenSessionAsync(request, "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
Assert.Equal($"client-{session.SessionId}", session.ClientCorrelationId);
|
||
}
|
||
|
||
/// <summary>Verifies that invoking a command on a ready session forwards the command to the worker.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenSessionReady_ForwardsCommandToWorker()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
WorkerCommandReply reply = await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal(1, workerClient.InvokeCount);
|
||
Assert.Equal(MxCommandKind.Ping, reply.Reply.Kind);
|
||
}
|
||
|
||
/// <summary>Verifies that invoking a command on a ready session refreshes its lease expiry.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenSessionReady_RefreshesLease()
|
||
{
|
||
GatewaySession session = new(
|
||
"session-lease-refresh",
|
||
"mxaccess",
|
||
"mxaccess-gateway-1-session-lease-refresh",
|
||
"nonce",
|
||
"client-1",
|
||
null,
|
||
"test-session",
|
||
"client-correlation-1",
|
||
TimeSpan.FromSeconds(30),
|
||
TimeSpan.FromSeconds(5),
|
||
TimeSpan.FromSeconds(5),
|
||
TimeSpan.FromMinutes(30),
|
||
DateTimeOffset.UtcNow - TimeSpan.FromHours(1));
|
||
session.AttachWorkerClient(new FakeWorkerClient());
|
||
session.MarkReady();
|
||
DateTimeOffset? initialLease = session.LeaseExpiresAt;
|
||
|
||
await session.InvokeAsync(CreateCommand(MxCommandKind.Ping), CancellationToken.None);
|
||
|
||
Assert.True(session.LeaseExpiresAt > initialLease);
|
||
Assert.True(session.LeaseExpiresAt > DateTimeOffset.UtcNow);
|
||
}
|
||
|
||
/// <summary>Verifies that gateway session subscribe bulk forwards one bulk command and returns results.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task GatewaySessionSubscribeBulkAsync_ForwardsOneBulkCommandAndReturnsResults()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
InvokeReply = new WorkerCommandReply
|
||
{
|
||
Reply = new MxCommandReply
|
||
{
|
||
SessionId = "session-1",
|
||
CorrelationId = "correlation-1",
|
||
Kind = MxCommandKind.SubscribeBulk,
|
||
ProtocolStatus = new ProtocolStatus { Code = ProtocolStatusCode.Ok },
|
||
SubscribeBulk = new BulkSubscribeReply
|
||
{
|
||
Results =
|
||
{
|
||
new SubscribeResult
|
||
{
|
||
ServerHandle = 12,
|
||
TagAddress = "Galaxy.Tag.Value",
|
||
ItemHandle = 512,
|
||
WasSuccessful = true,
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
IReadOnlyList<SubscribeResult> results = await session.SubscribeBulkAsync(
|
||
12,
|
||
["Galaxy.Tag.Value"],
|
||
CancellationToken.None);
|
||
|
||
SubscribeResult result = Assert.Single(results);
|
||
Assert.Equal(512, result.ItemHandle);
|
||
Assert.Equal(1, workerClient.InvokeCount);
|
||
Assert.Equal(MxCommandKind.SubscribeBulk, workerClient.LastCommand?.Command.Kind);
|
||
Assert.Equal(["Galaxy.Tag.Value"], workerClient.LastCommand?.Command.SubscribeBulk.TagAddresses);
|
||
}
|
||
|
||
/// <summary>Verifies that gateway session write bulk forwards one bulk command and returns results.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task GatewaySessionWriteBulkAsync_ForwardsOneBulkCommandAndReturnsResults()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
InvokeReply = new WorkerCommandReply
|
||
{
|
||
Reply = new MxCommandReply
|
||
{
|
||
SessionId = "session-1",
|
||
CorrelationId = "correlation-1",
|
||
Kind = MxCommandKind.WriteBulk,
|
||
ProtocolStatus = new ProtocolStatus { Code = ProtocolStatusCode.Ok },
|
||
WriteBulk = new BulkWriteReply
|
||
{
|
||
Results =
|
||
{
|
||
new BulkWriteResult
|
||
{
|
||
ServerHandle = 12,
|
||
ItemHandle = 901,
|
||
WasSuccessful = true,
|
||
},
|
||
new BulkWriteResult
|
||
{
|
||
ServerHandle = 12,
|
||
ItemHandle = 902,
|
||
WasSuccessful = false,
|
||
ErrorMessage = "MXAccess invalid handle",
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
IReadOnlyList<BulkWriteResult> results = await session.WriteBulkAsync(
|
||
12,
|
||
new[]
|
||
{
|
||
new WriteBulkEntry
|
||
{
|
||
ItemHandle = 901,
|
||
UserId = 5,
|
||
Value = new MxValue { DataType = MxDataType.Integer, Int32Value = 11 },
|
||
},
|
||
new WriteBulkEntry
|
||
{
|
||
ItemHandle = 902,
|
||
UserId = 5,
|
||
Value = new MxValue { DataType = MxDataType.Integer, Int32Value = 22 },
|
||
},
|
||
},
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal(2, results.Count);
|
||
Assert.True(results[0].WasSuccessful);
|
||
Assert.False(results[1].WasSuccessful);
|
||
Assert.Equal(MxCommandKind.WriteBulk, workerClient.LastCommand?.Command.Kind);
|
||
Assert.Equal(2, workerClient.LastCommand?.Command.WriteBulk.Entries.Count);
|
||
}
|
||
|
||
/// <summary>Verifies that gateway session read bulk forwards one bulk command and returns results.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task GatewaySessionReadBulkAsync_ForwardsOneBulkCommandAndReturnsResults()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
InvokeReply = new WorkerCommandReply
|
||
{
|
||
Reply = new MxCommandReply
|
||
{
|
||
SessionId = "session-1",
|
||
CorrelationId = "correlation-1",
|
||
Kind = MxCommandKind.ReadBulk,
|
||
ProtocolStatus = new ProtocolStatus { Code = ProtocolStatusCode.Ok },
|
||
ReadBulk = new BulkReadReply
|
||
{
|
||
Results =
|
||
{
|
||
new BulkReadResult
|
||
{
|
||
ServerHandle = 12,
|
||
TagAddress = "Galaxy.Tag.Value",
|
||
ItemHandle = 512,
|
||
WasSuccessful = true,
|
||
WasCached = true,
|
||
Value = new MxValue { DataType = MxDataType.Integer, Int32Value = 42 },
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
IReadOnlyList<BulkReadResult> results = await session.ReadBulkAsync(
|
||
12,
|
||
["Galaxy.Tag.Value"],
|
||
TimeSpan.FromMilliseconds(500),
|
||
CancellationToken.None);
|
||
|
||
BulkReadResult result = Assert.Single(results);
|
||
Assert.True(result.WasSuccessful);
|
||
Assert.True(result.WasCached);
|
||
Assert.Equal(42, result.Value.Int32Value);
|
||
Assert.Equal(MxCommandKind.ReadBulk, workerClient.LastCommand?.Command.Kind);
|
||
Assert.Equal(["Galaxy.Tag.Value"], workerClient.LastCommand?.Command.ReadBulk.TagAddresses);
|
||
Assert.Equal(500u, workerClient.LastCommand?.Command.ReadBulk.TimeoutMs);
|
||
}
|
||
|
||
/// <summary>Verifies that invoking a command on a faulted session rejects the command.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenSessionFaulted_RejectsCommand()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
session.MarkFaulted("test fault");
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
|
||
Assert.Equal(0, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// When the gateway-side <c>SessionState</c> is
|
||
/// <c>Ready</c> but the worker client's own state is not, the diagnostic
|
||
/// must surface both states so the mismatch is actionable instead of
|
||
/// producing a self-contradictory "Session ... is not ready. Current
|
||
/// state is Ready." message.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenWorkerNotReadyButSessionReady_DiagnosticIncludesBothStates()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
// Force a state mismatch: session stays Ready, worker transitions out.
|
||
workerClient.State = WorkerClientState.Handshaking;
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
|
||
Assert.Contains("Session state is Ready", exception.Message);
|
||
Assert.Contains("worker state is Handshaking", exception.Message);
|
||
Assert.Equal(0, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// With the opt-in worker-ready wait enabled, a worker that is transiently
|
||
/// <c>Handshaking</c> but flips to <c>Ready</c> within the timeout window must let the
|
||
/// command through rather than fail fast.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenWorkerHandshakingThenReadyWithinTimeout_Succeeds()
|
||
{
|
||
FakeWorkerClient workerClient = new() { State = WorkerClientState.Handshaking };
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(workerReadyWaitTimeoutMs: 500));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
// Flip the worker to Ready shortly after the invoke starts waiting.
|
||
_ = Task.Run(async () =>
|
||
{
|
||
await Task.Delay(50, CancellationToken.None);
|
||
workerClient.State = WorkerClientState.Ready;
|
||
});
|
||
|
||
WorkerCommandReply reply = await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None);
|
||
|
||
Assert.NotNull(reply);
|
||
Assert.Equal(1, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// A terminal worker state (<c>Faulted</c>) must fail fast even with a positive
|
||
/// worker-ready wait timeout, surfacing both states without burning the timeout.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenWorkerFaulted_FailsFastWithBothStates()
|
||
{
|
||
// Use a deliberately large ready-wait timeout so fail-fast is unambiguous: a terminal
|
||
// worker must surface immediately instead of burning it. The timing assertion below is
|
||
// anchored to a small fraction of this timeout (not an absolute ~100ms wall-clock bound,
|
||
// which flaked under CI load) so a regression that waited out the timeout is caught
|
||
// while ordinary scheduling jitter never trips it.
|
||
const int readyWaitTimeoutMs = 5000;
|
||
FakeWorkerClient workerClient = new() { State = WorkerClientState.Faulted };
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(workerReadyWaitTimeoutMs: readyWaitTimeoutMs));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
Stopwatch stopwatch = Stopwatch.StartNew();
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None));
|
||
stopwatch.Stop();
|
||
|
||
Assert.True(
|
||
stopwatch.ElapsedMilliseconds < readyWaitTimeoutMs / 3,
|
||
$"Expected fail-fast well under the {readyWaitTimeoutMs}ms ready-wait timeout but took {stopwatch.ElapsedMilliseconds}ms.");
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
|
||
Assert.Contains("Session state is Ready", exception.Message);
|
||
Assert.Contains("worker state is Faulted", exception.Message);
|
||
Assert.Equal(0, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// When the worker stays transiently not-ready for the whole (small) timeout window,
|
||
/// the invoke fails after roughly the timeout with both states surfaced.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenTimeoutElapsesStillNotReady_FailsWithBothStates()
|
||
{
|
||
FakeWorkerClient workerClient = new() { State = WorkerClientState.Handshaking };
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(workerReadyWaitTimeoutMs: 100));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
Stopwatch stopwatch = Stopwatch.StartNew();
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None));
|
||
stopwatch.Stop();
|
||
|
||
Assert.True(stopwatch.ElapsedMilliseconds >= 90, $"Expected the wait to span the timeout but took only {stopwatch.ElapsedMilliseconds}ms.");
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
|
||
Assert.Contains("Session state is Ready", exception.Message);
|
||
Assert.Contains("worker state is Handshaking", exception.Message);
|
||
Assert.Equal(0, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// Pins the default (timeout == 0) behavior: a transiently <c>Handshaking</c> worker
|
||
/// fails fast immediately, byte-for-byte like the original fail-fast path.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task InvokeAsync_WhenTimeoutZero_FailsFastUnchanged()
|
||
{
|
||
FakeWorkerClient workerClient = new() { State = WorkerClientState.Handshaking };
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(workerReadyWaitTimeoutMs: 0));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
// A zero timeout has no wait window to burn, so there is no wall-clock timing to assert
|
||
// (the previous absolute ~100ms bound only measured host load and flaked under CI). The
|
||
// error code, the byte-for-byte both-states message, and InvokeCount == 0 fully pin the
|
||
// immediate fail-fast — a regression that started polling would still surface the same
|
||
// outcome, not a timing difference worth a flaky assertion.
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.InvokeAsync(
|
||
session.SessionId,
|
||
CreateCommand(MxCommandKind.Ping),
|
||
CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
|
||
Assert.Contains("Session state is Ready", exception.Message);
|
||
Assert.Contains("worker state is Handshaking", exception.Message);
|
||
Assert.Equal(0, workerClient.InvokeCount);
|
||
}
|
||
|
||
/// <summary>Verifies that closing a session removes it from the registry.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseSessionAsync_RemovesClosedSession()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient), metrics: metrics);
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
SessionCloseResult firstClose = await manager.CloseSessionAsync(session.SessionId, CancellationToken.None);
|
||
SessionManagerException secondClose = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.CloseSessionAsync(session.SessionId, CancellationToken.None));
|
||
|
||
Assert.False(firstClose.AlreadyClosed);
|
||
Assert.Equal(SessionState.Closed, firstClose.FinalState);
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotFound, secondClose.ErrorCode);
|
||
Assert.Equal(1, workerClient.ShutdownCount);
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
}
|
||
|
||
/// <summary>Verifies that closing a session kills the worker when shutdown fails.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseSessionAsync_WhenWorkerShutdownFails_KillsWorker()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
ShutdownException = new WorkerClientException(
|
||
WorkerClientErrorCode.ShutdownTimeout,
|
||
"Worker shutdown timed out."),
|
||
};
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.CloseSessionAsync(session.SessionId, CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.CloseFailed, exception.ErrorCode);
|
||
Assert.Equal(1, workerClient.ShutdownCount);
|
||
Assert.Equal(1, workerClient.KillCount);
|
||
}
|
||
|
||
/// <summary>Verifies that when worker shutdown fails, the session is removed and the slot is released.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseSessionAsync_WhenWorkerShutdownFails_RemovesSessionAndReleasesSlot()
|
||
{
|
||
FakeWorkerClient failingWorkerClient = new()
|
||
{
|
||
ShutdownException = new WorkerClientException(
|
||
WorkerClientErrorCode.ShutdownTimeout,
|
||
"Worker shutdown timed out."),
|
||
};
|
||
FakeWorkerClient replacementWorkerClient = new();
|
||
SessionRegistry registry = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new QueueingSessionWorkerClientFactory(failingWorkerClient, replacementWorkerClient),
|
||
registry,
|
||
metrics,
|
||
CreateOptions(maxSessions: 1));
|
||
GatewaySession firstSession = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-1",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
metrics.EventReceived(firstSession.SessionId, MxEventFamily.OnDataChange.ToString());
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.CloseSessionAsync(firstSession.SessionId, CancellationToken.None));
|
||
GatewaySession secondSession = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-2",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal(SessionManagerErrorCode.CloseFailed, exception.ErrorCode);
|
||
Assert.False(manager.TryGetSession(firstSession.SessionId, out _));
|
||
Assert.True(manager.TryGetSession(secondSession.SessionId, out _));
|
||
Assert.Equal(1, registry.Count);
|
||
Assert.Equal(1, failingWorkerClient.KillCount);
|
||
Assert.Equal(1, failingWorkerClient.DisposeCount);
|
||
GatewayMetricsSnapshot snapshot = metrics.GetSnapshot();
|
||
Assert.Equal(1, snapshot.SessionsClosed);
|
||
Assert.False(snapshot.EventsBySession.ContainsKey(firstSession.SessionId));
|
||
Assert.Equal(1, snapshot.OpenSessions);
|
||
}
|
||
|
||
/// <summary>Verifies that when the second close is canceled, the session is not removed if owned by the first close.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseSessionAsync_WhenSecondCloseIsCanceled_DoesNotRemoveSessionOwnedByFirstClose()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
BlockShutdown = true,
|
||
};
|
||
SessionRegistry registry = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
registry,
|
||
metrics,
|
||
CreateOptions(maxSessions: 1));
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-1",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Task<SessionCloseResult> firstClose = manager.CloseSessionAsync(session.SessionId, CancellationToken.None);
|
||
await workerClient.WaitForShutdownStartAsync();
|
||
using CancellationTokenSource secondCloseCancellation = new();
|
||
Task<SessionCloseResult> secondClose = manager.CloseSessionAsync(
|
||
session.SessionId,
|
||
secondCloseCancellation.Token);
|
||
|
||
await secondCloseCancellation.CancelAsync();
|
||
|
||
await Assert.ThrowsAnyAsync<OperationCanceledException>(
|
||
async () => await secondClose);
|
||
Assert.True(manager.TryGetSession(session.SessionId, out _));
|
||
Assert.Equal(1, registry.Count);
|
||
Assert.Equal(0, workerClient.DisposeCount);
|
||
Assert.Equal(0, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(1, metrics.GetSnapshot().OpenSessions);
|
||
|
||
workerClient.ReleaseShutdown();
|
||
SessionCloseResult closeResult = await firstClose;
|
||
|
||
Assert.Equal(SessionState.Closed, closeResult.FinalState);
|
||
Assert.False(manager.TryGetSession(session.SessionId, out _));
|
||
Assert.Equal(0, registry.Count);
|
||
Assert.Equal(1, workerClient.DisposeCount);
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
}
|
||
|
||
/// <summary>
|
||
/// Verifies that killing a worker removes the session from the registry without calling shutdown.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task KillWorkerAsync_KillsWorkerAndRemovesSession()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient), metrics: metrics);
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
SessionCloseResult result = await manager.KillWorkerAsync(session.SessionId, "test-kill", CancellationToken.None);
|
||
|
||
Assert.False(result.AlreadyClosed);
|
||
Assert.Equal(SessionState.Closed, result.FinalState);
|
||
Assert.Equal(1, workerClient.KillCount);
|
||
Assert.Equal("test-kill", workerClient.LastKillReason);
|
||
Assert.Equal(0, workerClient.ShutdownCount);
|
||
Assert.False(manager.TryGetSession(session.SessionId, out _));
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
}
|
||
|
||
/// <summary>
|
||
/// <see cref="SessionManager.KillWorkerAsync"/> guards its <c>reason</c> argument with
|
||
/// <see cref="ArgumentException.ThrowIfNullOrWhiteSpace"/>. A blank or whitespace reason must throw
|
||
/// <see cref="ArgumentException"/> before any session lookup or worker call runs.
|
||
/// </summary>
|
||
/// <param name="blankReason">A blank or whitespace reason string.</param>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Theory]
|
||
[InlineData("")]
|
||
[InlineData(" ")]
|
||
[InlineData("\t")]
|
||
public async Task KillWorkerAsync_WithBlankReason_ThrowsArgumentException(string blankReason)
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
await Assert.ThrowsAsync<ArgumentException>(
|
||
async () => await manager.KillWorkerAsync(session.SessionId, blankReason, CancellationToken.None));
|
||
|
||
Assert.Equal(0, workerClient.KillCount);
|
||
Assert.True(manager.TryGetSession(session.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>
|
||
/// <see cref="ArgumentException.ThrowIfNullOrWhiteSpace"/> also rejects null.
|
||
/// <see cref="Theory"/> with <see cref="InlineDataAttribute"/> cannot carry <c>null</c> for a
|
||
/// non-nullable string parameter on .NET 10, so the null case is its own fact.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task KillWorkerAsync_WithNullReason_ThrowsArgumentNullException()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
await Assert.ThrowsAsync<ArgumentNullException>(
|
||
async () => await manager.KillWorkerAsync(session.SessionId, null!, CancellationToken.None));
|
||
|
||
Assert.Equal(0, workerClient.KillCount);
|
||
Assert.True(manager.TryGetSession(session.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>Verifies that killing the worker for an unknown session raises SessionNotFound.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task KillWorkerAsync_WhenSessionMissing_ThrowsSessionNotFound()
|
||
{
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(new FakeWorkerClient()));
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.KillWorkerAsync("session-missing", "test-kill", CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.SessionNotFound, exception.ErrorCode);
|
||
}
|
||
|
||
/// <summary>
|
||
/// When <c>session.KillWorker</c> throws, the catch path must still
|
||
/// decrement <c>mxgateway.sessions.open</c>. Without the fix the gauge leaks one open session per failed kill.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task KillWorkerAsync_WhenSessionKillThrows_DecrementsOpenSessionGauge()
|
||
{
|
||
FakeWorkerClient workerClient = new()
|
||
{
|
||
KillException = new InvalidOperationException("worker kill failed"),
|
||
};
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
metrics: metrics);
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-1",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal(1, metrics.GetSnapshot().OpenSessions);
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.KillWorkerAsync(session.SessionId, "test-kill", CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.CloseFailed, exception.ErrorCode);
|
||
Assert.False(manager.TryGetSession(session.SessionId, out _));
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
Assert.True(metrics.GetSnapshot().Faults > 0);
|
||
}
|
||
|
||
/// <summary>
|
||
/// Concurrent kills on the same session must not
|
||
/// double-increment <c>mxgateway.sessions.closed</c>. The first kill wins, the second
|
||
/// observes <c>wasClosed == true</c> (or a missing session after removal) and short-circuits.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task KillWorkerAsync_ConcurrentCallsOnSameSession_CountClosedExactlyOnce()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
metrics: metrics);
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-1",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Task<SessionCloseResult> first = manager.KillWorkerAsync(session.SessionId, "kill-a", CancellationToken.None);
|
||
Task<SessionCloseResult> second = Task.Run(async () =>
|
||
{
|
||
try
|
||
{
|
||
return await manager.KillWorkerAsync(session.SessionId, "kill-b", CancellationToken.None);
|
||
}
|
||
catch (SessionManagerException missing) when (missing.ErrorCode == SessionManagerErrorCode.SessionNotFound)
|
||
{
|
||
return new SessionCloseResult(session.SessionId, SessionState.Closed, AlreadyClosed: true);
|
||
}
|
||
});
|
||
|
||
await Task.WhenAll(first, second);
|
||
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
Assert.False(manager.TryGetSession(session.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>
|
||
/// <c>ShutdownAsync</c>'s graceful-close fallback (which calls
|
||
/// <c>KillWorker</c> + <c>RemoveSessionAsync</c> when <c>CloseSessionCoreAsync</c> throws)
|
||
/// must still account a successful close: both the open-session gauge must drop to zero AND
|
||
/// the <c>mxgateway.sessions.closed</c> counter must increment. Without the fix, the
|
||
/// graceful-close failure path under-counts the closed counter.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task ShutdownAsync_WhenSessionCloseThrows_StillDecrementsOpenSessionGaugeAndIncrementsClosedCounter()
|
||
{
|
||
FakeWorkerClient throwingClient = new()
|
||
{
|
||
ShutdownException = new InvalidOperationException("worker shutdown failed"),
|
||
};
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(throwingClient),
|
||
metrics: metrics);
|
||
GatewaySession session = await manager.OpenSessionAsync(
|
||
CreateOpenRequest(),
|
||
"client-1",
|
||
ownerKeyId: null,
|
||
CancellationToken.None);
|
||
|
||
Assert.Equal(1, metrics.GetSnapshot().OpenSessions);
|
||
|
||
await manager.ShutdownAsync(CancellationToken.None);
|
||
|
||
// After shutdown, regardless of whether the graceful close path or the kill fallback ran,
|
||
// the open-session gauge must be zero and the closed counter must be incremented.
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
Assert.Equal(1, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.False(manager.TryGetSession(session.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>Verifies that when worker creation fails, the session is removed from the registry.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task OpenSessionAsync_WhenWorkerCreationFails_RemovesSessionFromRegistry()
|
||
{
|
||
SessionRegistry registry = new();
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(
|
||
new FailingSessionWorkerClientFactory(),
|
||
registry,
|
||
metrics);
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.OpenFailed, exception.ErrorCode);
|
||
Assert.Equal(0, registry.Count);
|
||
Assert.Equal(0, metrics.GetSnapshot().SessionsOpened);
|
||
Assert.Equal(1, metrics.GetSnapshot().Faults);
|
||
}
|
||
|
||
/// <summary>Verifies that closing expired leases only closes expired sessions.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_ClosesExpiredSessionsOnly()
|
||
{
|
||
FakeWorkerClient expiredClient = new();
|
||
FakeWorkerClient activeClient = new();
|
||
QueueingSessionWorkerClientFactory factory = new(expiredClient, activeClient);
|
||
SessionManager manager = CreateManager(factory);
|
||
GatewaySession expiredSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession activeSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
DateTimeOffset now = DateTimeOffset.UtcNow;
|
||
expiredSession.ExtendLease(now.AddSeconds(-1));
|
||
activeSession.ExtendLease(now.AddMinutes(5));
|
||
|
||
int closedCount = await manager.CloseExpiredLeasesAsync(now, CancellationToken.None);
|
||
|
||
Assert.Equal(1, closedCount);
|
||
Assert.Equal(SessionState.Closed, expiredSession.State);
|
||
Assert.Equal(SessionState.Ready, activeSession.State);
|
||
Assert.Equal(1, expiredClient.ShutdownCount);
|
||
Assert.Equal(0, activeClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>Verifies that an expired-lease sweep leaves a session with an active event subscriber open.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_DoesNotCloseActiveEventSubscriber()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
SessionManager manager = CreateManager(new FakeSessionWorkerClientFactory(workerClient));
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
DateTimeOffset now = DateTimeOffset.UtcNow;
|
||
session.ExtendLease(now.AddSeconds(-1));
|
||
using IDisposable eventSubscriber = session.AttachEventSubscriber(maxSubscribers: 1);
|
||
|
||
int closedCount = await manager.CloseExpiredLeasesAsync(now, CancellationToken.None);
|
||
|
||
Assert.Equal(0, closedCount);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
Assert.Equal(0, workerClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// With detach-grace enabled, a session whose last external subscriber dropped
|
||
/// and whose detach-grace window has elapsed is closed by the lease sweep exactly like an
|
||
/// expired-lease session — even though its normal lease is still far in the future.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_ClosesSessionWhoseDetachGraceWindowExpired()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
FakeTimeProvider clock = new(DateTimeOffset.UtcNow);
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(defaultLeaseSeconds: 1800, detachGraceSeconds: 30),
|
||
timeProvider: clock);
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
// Attach and drop an external subscriber to enter detach-grace. The normal lease is
|
||
// still 30 minutes out, so only the detach-grace window can close this session.
|
||
IDisposable subscriber = session.AttachEventSubscriber(maxSubscribers: 1);
|
||
subscriber.Dispose();
|
||
Assert.NotNull(session.DetachedAtUtc);
|
||
|
||
// Before the window elapses: not closed.
|
||
clock.Advance(TimeSpan.FromSeconds(29));
|
||
int closedBefore = await manager.CloseExpiredLeasesAsync(clock.GetUtcNow(), CancellationToken.None);
|
||
Assert.Equal(0, closedBefore);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
|
||
// After the window elapses: the sweep closes it.
|
||
clock.Advance(TimeSpan.FromSeconds(1));
|
||
int closedAfter = await manager.CloseExpiredLeasesAsync(clock.GetUtcNow(), CancellationToken.None);
|
||
Assert.Equal(1, closedAfter);
|
||
Assert.Equal(SessionState.Closed, session.State);
|
||
Assert.Equal(1, workerClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// TOCTOU race: a session whose detach-grace window has expired but that
|
||
/// reattaches an external subscriber before the sweeper calls CloseSessionCoreAsync is
|
||
/// NOT closed — it remains Ready and usable. This validates that TryBeginCloseIfExpired
|
||
/// re-checks eligibility atomically so a reconnect that wins the race cancels the close.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_DoesNotCloseSessionThatReattachedBeforeSweepCloses()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
FakeTimeProvider clock = new(DateTimeOffset.UtcNow);
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(defaultLeaseSeconds: 1800, detachGraceSeconds: 30),
|
||
timeProvider: clock);
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
|
||
// Attach and drop an external subscriber so the session enters detach-grace.
|
||
IDisposable firstSubscriber = session.AttachEventSubscriber(maxSubscribers: 1);
|
||
firstSubscriber.Dispose();
|
||
Assert.NotNull(session.DetachedAtUtc);
|
||
|
||
// Advance past the grace window so IsDetachGraceExpired returns true.
|
||
clock.Advance(TimeSpan.FromSeconds(31));
|
||
DateTimeOffset sweepTime = clock.GetUtcNow();
|
||
|
||
// Simulate a client reattaching before the sweep actually closes the session.
|
||
// The reattach clears _detachedAtUtc and increments _activeEventSubscriberCount,
|
||
// so TryBeginCloseIfExpired will see neither condition as met and decline.
|
||
using IDisposable reconnectedSubscriber = session.AttachEventSubscriber(maxSubscribers: 1);
|
||
Assert.Null(session.DetachedAtUtc);
|
||
|
||
// The sweep runs with the timestamp that was past the grace window, but since the
|
||
// subscriber has reattached, the session must NOT be closed.
|
||
int closedCount = await manager.CloseExpiredLeasesAsync(sweepTime, CancellationToken.None);
|
||
|
||
Assert.Equal(0, closedCount);
|
||
Assert.Equal(SessionState.Ready, session.State);
|
||
Assert.Equal(0, workerClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// A faulted session is reaped by the lease sweep (default <c>FaultedGraceSeconds=0</c>)
|
||
/// even though its normal lease is still far in the future, tearing down its worker and
|
||
/// freeing the slot, while a healthy leased session in the same manager is untouched.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_ReapsFaultedSession()
|
||
{
|
||
FakeWorkerClient faultedClient = new();
|
||
FakeWorkerClient healthyClient = new();
|
||
QueueingSessionWorkerClientFactory factory = new(faultedClient, healthyClient);
|
||
SessionManager manager = CreateManager(factory);
|
||
GatewaySession faultedSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession healthySession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
DateTimeOffset now = DateTimeOffset.UtcNow;
|
||
|
||
// Both leases are far in the future, so only the fault can reap the faulted session.
|
||
faultedSession.ExtendLease(now.AddMinutes(30));
|
||
healthySession.ExtendLease(now.AddMinutes(30));
|
||
faultedSession.MarkFaulted("test fault");
|
||
|
||
int closedCount = await manager.CloseExpiredLeasesAsync(now, CancellationToken.None);
|
||
|
||
Assert.Equal(1, closedCount);
|
||
Assert.Equal(SessionState.Closed, faultedSession.State);
|
||
Assert.Equal(1, faultedClient.ShutdownCount);
|
||
Assert.Equal(SessionState.Ready, healthySession.State);
|
||
Assert.Equal(0, healthyClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// With a positive <c>FaultedGraceSeconds</c>, a faulted session stays observable until
|
||
/// the grace window elapses, then the next sweep reaps it.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_RespectsFaultedGraceWindow()
|
||
{
|
||
FakeWorkerClient workerClient = new();
|
||
FakeTimeProvider clock = new(DateTimeOffset.UtcNow);
|
||
SessionManager manager = CreateManager(
|
||
new FakeSessionWorkerClientFactory(workerClient),
|
||
options: CreateOptions(defaultLeaseSeconds: 1800, faultedGraceSeconds: 30),
|
||
timeProvider: clock);
|
||
GatewaySession session = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
session.MarkFaulted("test fault");
|
||
|
||
// Before the grace window elapses: the faulted session is retained (still observable).
|
||
clock.Advance(TimeSpan.FromSeconds(29));
|
||
int closedBefore = await manager.CloseExpiredLeasesAsync(clock.GetUtcNow(), CancellationToken.None);
|
||
Assert.Equal(0, closedBefore);
|
||
Assert.Equal(SessionState.Faulted, session.State);
|
||
|
||
// After the grace window elapses: the sweep reaps it.
|
||
clock.Advance(TimeSpan.FromSeconds(1));
|
||
int closedAfter = await manager.CloseExpiredLeasesAsync(clock.GetUtcNow(), CancellationToken.None);
|
||
Assert.Equal(1, closedAfter);
|
||
Assert.Equal(SessionState.Closed, session.State);
|
||
Assert.Equal(1, workerClient.ShutdownCount);
|
||
}
|
||
|
||
/// <summary>
|
||
/// A sweep pass tears the selected sessions down concurrently rather than one worker
|
||
/// shutdown after another: with a mass expiry, a few hung workers would otherwise
|
||
/// serialize reaping (each close is bounded by <c>Worker:ShutdownTimeoutSeconds</c>) and
|
||
/// starve session slots. Overlap is asserted by counting concurrent entries into the fake
|
||
/// worker's shutdown rather than by wall clock, which is sturdier on a loaded box.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_ClosesExpiredSessionsConcurrently()
|
||
{
|
||
ShutdownConcurrencyProbe probe = new(expectedConcurrency: 2);
|
||
FakeWorkerClient firstClient = new() { ShutdownConcurrencyProbe = probe };
|
||
FakeWorkerClient secondClient = new() { ShutdownConcurrencyProbe = probe };
|
||
SessionManager manager = CreateManager(new QueueingSessionWorkerClientFactory(firstClient, secondClient));
|
||
GatewaySession firstSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession secondSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
DateTimeOffset now = DateTimeOffset.UtcNow;
|
||
firstSession.ExtendLease(now.AddSeconds(-1));
|
||
secondSession.ExtendLease(now.AddSeconds(-1));
|
||
|
||
int closedCount = await manager.CloseExpiredLeasesAsync(now, CancellationToken.None);
|
||
|
||
Assert.Equal(2, closedCount);
|
||
Assert.Equal(SessionState.Closed, firstSession.State);
|
||
Assert.Equal(SessionState.Closed, secondSession.State);
|
||
Assert.Equal(2, probe.MaxObservedConcurrency);
|
||
}
|
||
|
||
/// <summary>
|
||
/// Host stop drains sessions concurrently: 50 sessions at a worst-case 10 s shutdown each
|
||
/// would exceed any host stop-timeout if drained one at a time, leaving the tail to the
|
||
/// orphan killer.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task ShutdownAsync_ClosesSessionsConcurrently()
|
||
{
|
||
ShutdownConcurrencyProbe probe = new(expectedConcurrency: 2);
|
||
FakeWorkerClient firstClient = new() { ShutdownConcurrencyProbe = probe };
|
||
FakeWorkerClient secondClient = new() { ShutdownConcurrencyProbe = probe };
|
||
SessionManager manager = CreateManager(new QueueingSessionWorkerClientFactory(firstClient, secondClient));
|
||
GatewaySession firstSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession secondSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
|
||
await manager.ShutdownAsync(CancellationToken.None);
|
||
|
||
Assert.Equal(SessionState.Closed, firstSession.State);
|
||
Assert.Equal(SessionState.Closed, secondSession.State);
|
||
Assert.Equal(2, probe.MaxObservedConcurrency);
|
||
}
|
||
|
||
/// <summary>
|
||
/// A close that throws must not abandon the rest of the selected set: the sweep still
|
||
/// tears the healthy expired session down, and the failure still surfaces to the lease
|
||
/// monitor (which logs it) exactly as the sequential loop did.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task CloseExpiredLeasesAsync_WhenOneCloseFails_StillClosesRemainingSessionsAndRethrows()
|
||
{
|
||
FakeWorkerClient failingClient = new()
|
||
{
|
||
ShutdownException = new InvalidOperationException("worker shutdown failed"),
|
||
KillException = new InvalidOperationException("worker kill failed"),
|
||
};
|
||
FakeWorkerClient healthyClient = new();
|
||
SessionManager manager = CreateManager(new QueueingSessionWorkerClientFactory(failingClient, healthyClient));
|
||
GatewaySession failingSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession healthySession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
DateTimeOffset now = DateTimeOffset.UtcNow;
|
||
failingSession.ExtendLease(now.AddSeconds(-1));
|
||
healthySession.ExtendLease(now.AddSeconds(-1));
|
||
|
||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||
async () => await manager.CloseExpiredLeasesAsync(now, CancellationToken.None));
|
||
|
||
Assert.Equal(SessionManagerErrorCode.CloseFailed, exception.ErrorCode);
|
||
Assert.Equal(1, healthyClient.ShutdownCount);
|
||
Assert.Equal(SessionState.Closed, healthySession.State);
|
||
Assert.False(manager.TryGetSession(healthySession.SessionId, out _));
|
||
Assert.False(manager.TryGetSession(failingSession.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>
|
||
/// A drain whose token is already cancelled (the host stop deadline elapsed) must still
|
||
/// kill every worker rather than skip the teardown: an unkilled worker is a leaked x86
|
||
/// process, and a restarted gateway terminates orphans instead of reattaching to them. This
|
||
/// pins both halves of the fix — the parallel loop is not bound to the caller's token, and
|
||
/// the kill fallback does not run on it.
|
||
/// </summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task ShutdownAsync_WhenCancelledBeforeDraining_StillKillsEveryWorker()
|
||
{
|
||
FakeWorkerClient firstClient = new();
|
||
FakeWorkerClient secondClient = new();
|
||
SessionManager manager = CreateManager(new QueueingSessionWorkerClientFactory(firstClient, secondClient));
|
||
GatewaySession firstSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession secondSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
using CancellationTokenSource cancellation = new();
|
||
await cancellation.CancelAsync();
|
||
|
||
await manager.ShutdownAsync(cancellation.Token);
|
||
|
||
Assert.Equal(1, firstClient.KillCount);
|
||
Assert.Equal(1, secondClient.KillCount);
|
||
Assert.False(manager.TryGetSession(firstSession.SessionId, out _));
|
||
Assert.False(manager.TryGetSession(secondSession.SessionId, out _));
|
||
}
|
||
|
||
/// <summary>Verifies that shutdown closes all registered sessions.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
[Fact]
|
||
public async Task ShutdownAsync_ClosesAllRegisteredSessions()
|
||
{
|
||
FakeWorkerClient firstClient = new();
|
||
FakeWorkerClient secondClient = new();
|
||
QueueingSessionWorkerClientFactory factory = new(firstClient, secondClient);
|
||
using GatewayMetrics metrics = new();
|
||
SessionManager manager = CreateManager(factory, metrics: metrics);
|
||
GatewaySession firstSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-1", ownerKeyId: null, CancellationToken.None);
|
||
GatewaySession secondSession = await manager.OpenSessionAsync(CreateOpenRequest(), "client-2", ownerKeyId: null, CancellationToken.None);
|
||
|
||
await manager.ShutdownAsync(CancellationToken.None);
|
||
|
||
Assert.Equal(SessionState.Closed, firstSession.State);
|
||
Assert.Equal(SessionState.Closed, secondSession.State);
|
||
Assert.Equal(1, firstClient.ShutdownCount);
|
||
Assert.Equal(1, secondClient.ShutdownCount);
|
||
Assert.Equal(2, metrics.GetSnapshot().SessionsClosed);
|
||
Assert.Equal(0, metrics.GetSnapshot().OpenSessions);
|
||
}
|
||
|
||
/// <summary>Creates a session manager for testing.</summary>
|
||
/// <param name="factory">Worker client factory.</param>
|
||
/// <param name="registry">Session registry; defaults to a new registry.</param>
|
||
/// <param name="metrics">Metrics collector; defaults to a new instance.</param>
|
||
/// <param name="options">Gateway options; defaults to test defaults.</param>
|
||
/// <returns>Configured session manager.</returns>
|
||
private static SessionManager CreateManager(
|
||
ISessionWorkerClientFactory factory,
|
||
ISessionRegistry? registry = null,
|
||
GatewayMetrics? metrics = null,
|
||
GatewayOptions? options = null,
|
||
TimeProvider? timeProvider = null)
|
||
{
|
||
return new SessionManager(
|
||
registry ?? new SessionRegistry(),
|
||
factory,
|
||
Options.Create(options ?? CreateOptions()),
|
||
metrics ?? new GatewayMetrics(),
|
||
timeProvider);
|
||
}
|
||
|
||
private static GatewayOptions CreateOptions(
|
||
int maxSessions = 64,
|
||
int defaultLeaseSeconds = 1800,
|
||
int detachGraceSeconds = 0,
|
||
int workerReadyWaitTimeoutMs = 0,
|
||
int faultedGraceSeconds = 0)
|
||
{
|
||
return new GatewayOptions
|
||
{
|
||
Sessions = new SessionOptions
|
||
{
|
||
DefaultCommandTimeoutSeconds = 30,
|
||
MaxSessions = maxSessions,
|
||
DefaultLeaseSeconds = defaultLeaseSeconds,
|
||
DetachGraceSeconds = detachGraceSeconds,
|
||
WorkerReadyWaitTimeoutMs = workerReadyWaitTimeoutMs,
|
||
FaultedGraceSeconds = faultedGraceSeconds,
|
||
},
|
||
Worker = new WorkerOptions
|
||
{
|
||
StartupTimeoutSeconds = 30,
|
||
ShutdownTimeoutSeconds = 10,
|
||
},
|
||
};
|
||
}
|
||
|
||
private static SessionOpenRequest CreateOpenRequest()
|
||
{
|
||
return new SessionOpenRequest(
|
||
RequestedBackend: null,
|
||
ClientSessionName: "test-session",
|
||
ClientCorrelationId: "client-correlation-1",
|
||
CommandTimeout: Duration.FromTimeSpan(TimeSpan.FromSeconds(5)));
|
||
}
|
||
|
||
private static WorkerCommand CreateCommand(MxCommandKind kind)
|
||
{
|
||
return new WorkerCommand
|
||
{
|
||
Command = new MxCommand
|
||
{
|
||
Kind = kind,
|
||
},
|
||
};
|
||
}
|
||
|
||
private sealed class FakeSessionWorkerClientFactory(IWorkerClient workerClient) : ISessionWorkerClientFactory
|
||
{
|
||
/// <summary>Gets the list of observed session states during worker creation.</summary>
|
||
public List<string> ObservedStates { get; } = [];
|
||
|
||
/// <summary>Gets or sets a value indicating whether to apply lifecycle transitions during worker creation.</summary>
|
||
public bool ApplyLifecycleTransitions { get; init; }
|
||
|
||
/// <inheritdoc />
|
||
public Task<IWorkerClient> CreateAsync(
|
||
GatewaySession session,
|
||
CancellationToken cancellationToken)
|
||
{
|
||
ObservedStates.Add(session.State.ToString());
|
||
if (ApplyLifecycleTransitions)
|
||
{
|
||
session.TransitionTo(SessionState.WaitingForPipe);
|
||
ObservedStates.Add(session.State.ToString());
|
||
session.TransitionTo(SessionState.Handshaking);
|
||
ObservedStates.Add(session.State.ToString());
|
||
session.TransitionTo(SessionState.InitializingWorker);
|
||
ObservedStates.Add(session.State.ToString());
|
||
}
|
||
|
||
return Task.FromResult(workerClient);
|
||
}
|
||
}
|
||
|
||
private sealed class QueueingSessionWorkerClientFactory : ISessionWorkerClientFactory
|
||
{
|
||
private readonly Queue<IWorkerClient> _workerClients;
|
||
|
||
/// <summary>Initializes a new instance of the <see cref="QueueingSessionWorkerClientFactory"/> class.</summary>
|
||
/// <param name="workerClients">Array of worker clients to queue.</param>
|
||
public QueueingSessionWorkerClientFactory(params IWorkerClient[] workerClients)
|
||
{
|
||
_workerClients = new Queue<IWorkerClient>(workerClients);
|
||
}
|
||
|
||
/// <inheritdoc />
|
||
public Task<IWorkerClient> CreateAsync(
|
||
GatewaySession session,
|
||
CancellationToken cancellationToken)
|
||
{
|
||
return Task.FromResult(_workerClients.Dequeue());
|
||
}
|
||
}
|
||
|
||
private sealed class FailingSessionWorkerClientFactory : ISessionWorkerClientFactory
|
||
{
|
||
/// <inheritdoc />
|
||
public Task<IWorkerClient> CreateAsync(
|
||
GatewaySession session,
|
||
CancellationToken cancellationToken)
|
||
{
|
||
throw new InvalidOperationException("worker startup failed");
|
||
}
|
||
}
|
||
|
||
private sealed class FakeWorkerClient : IWorkerClient
|
||
{
|
||
/// <inheritdoc />
|
||
public string SessionId { get; init; } = "session-1";
|
||
|
||
/// <inheritdoc />
|
||
public int? ProcessId { get; init; } = 1234;
|
||
|
||
/// <inheritdoc />
|
||
public WorkerClientState State { get; set; } = WorkerClientState.Ready;
|
||
|
||
/// <inheritdoc />
|
||
public DateTimeOffset LastHeartbeatAt { get; init; } = DateTimeOffset.UtcNow;
|
||
|
||
/// <summary>Gets the number of times invoke was called on the fake worker client.</summary>
|
||
public int InvokeCount { get; private set; }
|
||
|
||
/// <summary>Gets the number of times shutdown was called on the fake worker client.</summary>
|
||
public int ShutdownCount { get; private set; }
|
||
|
||
/// <summary>Gets the number of times kill was called on the fake worker client.</summary>
|
||
public int KillCount { get; private set; }
|
||
|
||
/// <summary>Gets the last reason argument observed by <see cref="Kill"/>.</summary>
|
||
public string? LastKillReason { get; private set; }
|
||
|
||
/// <summary>Gets the number of times dispose was called on the fake worker client.</summary>
|
||
public int DisposeCount { get; private set; }
|
||
|
||
/// <summary>Gets the exception to throw when shutdown is called, if any.</summary>
|
||
public Exception? ShutdownException { get; init; }
|
||
|
||
/// <summary>Gets the exception to throw when kill is called, if any.</summary>
|
||
public Exception? KillException { get; init; }
|
||
|
||
/// <summary>Gets a value indicating whether to block shutdown on the fake worker client.</summary>
|
||
public bool BlockShutdown { get; init; }
|
||
|
||
/// <summary>
|
||
/// Gets the rendezvous that records how many shutdowns overlap, shared by the fakes of
|
||
/// the sessions a single teardown pass closes. Null when the test does not measure
|
||
/// teardown concurrency.
|
||
/// </summary>
|
||
public ShutdownConcurrencyProbe? ShutdownConcurrencyProbe { get; init; }
|
||
|
||
/// <summary>Gets the last command invoked on the fake worker client.</summary>
|
||
public WorkerCommand? LastCommand { get; private set; }
|
||
|
||
/// <summary>Gets the reply to return for invoke calls on the fake worker client.</summary>
|
||
public WorkerCommandReply? InvokeReply { get; init; }
|
||
|
||
private TaskCompletionSource ShutdownStarted { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
|
||
|
||
private TaskCompletionSource ShutdownReleased { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
|
||
|
||
/// <inheritdoc />
|
||
public Task StartAsync(CancellationToken cancellationToken)
|
||
{
|
||
return Task.CompletedTask;
|
||
}
|
||
|
||
/// <inheritdoc />
|
||
public Task<WorkerCommandReply> InvokeAsync(
|
||
WorkerCommand command,
|
||
TimeSpan timeout,
|
||
CancellationToken cancellationToken)
|
||
{
|
||
InvokeCount++;
|
||
LastCommand = command;
|
||
if (InvokeReply is not null)
|
||
{
|
||
return Task.FromResult(InvokeReply);
|
||
}
|
||
|
||
MxCommandKind kind = command.Command?.Kind ?? MxCommandKind.Unspecified;
|
||
|
||
return Task.FromResult(new WorkerCommandReply
|
||
{
|
||
Reply = new MxCommandReply
|
||
{
|
||
SessionId = SessionId,
|
||
CorrelationId = "correlation-1",
|
||
Kind = kind,
|
||
},
|
||
});
|
||
}
|
||
|
||
/// <inheritdoc />
|
||
public async IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
|
||
[System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken)
|
||
{
|
||
await Task.CompletedTask;
|
||
yield break;
|
||
}
|
||
|
||
/// <inheritdoc />
|
||
public async Task ShutdownAsync(
|
||
TimeSpan timeout,
|
||
CancellationToken cancellationToken)
|
||
{
|
||
ShutdownCount++;
|
||
if (ShutdownException is not null)
|
||
{
|
||
throw ShutdownException;
|
||
}
|
||
|
||
if (ShutdownConcurrencyProbe is not null)
|
||
{
|
||
await ShutdownConcurrencyProbe.EnterAsync(cancellationToken);
|
||
}
|
||
|
||
if (BlockShutdown)
|
||
{
|
||
ShutdownStarted.TrySetResult();
|
||
await ShutdownReleased.Task.WaitAsync(cancellationToken);
|
||
}
|
||
|
||
State = WorkerClientState.Closed;
|
||
}
|
||
|
||
/// <inheritdoc />
|
||
public void Kill(string reason)
|
||
{
|
||
KillCount++;
|
||
LastKillReason = reason;
|
||
if (KillException is not null)
|
||
{
|
||
throw KillException;
|
||
}
|
||
|
||
State = WorkerClientState.Faulted;
|
||
}
|
||
|
||
/// <summary>Records that the fake worker client was disposed.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
public ValueTask DisposeAsync()
|
||
{
|
||
DisposeCount++;
|
||
return ValueTask.CompletedTask;
|
||
}
|
||
|
||
/// <summary>Waits for shutdown to start on the fake worker client.</summary>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
public Task WaitForShutdownStartAsync()
|
||
{
|
||
return ShutdownStarted.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
||
}
|
||
|
||
/// <summary>Releases the shutdown block on the fake worker client.</summary>
|
||
public void ReleaseShutdown()
|
||
{
|
||
ShutdownReleased.TrySetResult();
|
||
}
|
||
}
|
||
|
||
/// <summary>
|
||
/// Rendezvous that measures how many worker shutdowns a teardown pass runs at once. Each
|
||
/// entering shutdown records the in-flight count and waits until <paramref name="expectedConcurrency"/>
|
||
/// shutdowns are in flight, so a genuinely parallel teardown releases immediately while a
|
||
/// sequential one can only release on the bounded timeout — with a max observed concurrency
|
||
/// of one, which is the assertion that fails.
|
||
/// </summary>
|
||
/// <param name="expectedConcurrency">Number of overlapping shutdowns that releases the rendezvous.</param>
|
||
private sealed class ShutdownConcurrencyProbe(int expectedConcurrency)
|
||
{
|
||
private static readonly TimeSpan RendezvousTimeout = TimeSpan.FromSeconds(5);
|
||
|
||
private readonly TaskCompletionSource _reached = new(TaskCreationOptions.RunContinuationsAsynchronously);
|
||
private int _inFlight;
|
||
private int _maxInFlight;
|
||
|
||
/// <summary>Gets the highest number of shutdowns observed in flight at the same time.</summary>
|
||
public int MaxObservedConcurrency => Volatile.Read(ref _maxInFlight);
|
||
|
||
/// <summary>Enters the rendezvous for one worker shutdown and waits for the expected overlap.</summary>
|
||
/// <param name="cancellationToken">Token that abandons the wait.</param>
|
||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||
public async Task EnterAsync(CancellationToken cancellationToken)
|
||
{
|
||
int inFlight = Interlocked.Increment(ref _inFlight);
|
||
RecordMax(inFlight);
|
||
if (inFlight >= expectedConcurrency)
|
||
{
|
||
_reached.TrySetResult();
|
||
}
|
||
|
||
try
|
||
{
|
||
await _reached.Task.WaitAsync(RendezvousTimeout, cancellationToken);
|
||
}
|
||
catch (Exception exception) when (exception is TimeoutException or OperationCanceledException)
|
||
{
|
||
// Sequential teardown: the expected overlap never happens, so let the shutdown
|
||
// finish and let MaxObservedConcurrency report the (failing) truth. Cancellation is
|
||
// swallowed for the same reason — a cancelled rendezvous must not turn into a
|
||
// second, misleading failure on top of the concurrency assertion.
|
||
}
|
||
finally
|
||
{
|
||
Interlocked.Decrement(ref _inFlight);
|
||
}
|
||
}
|
||
|
||
private void RecordMax(int inFlight)
|
||
{
|
||
int observed = Volatile.Read(ref _maxInFlight);
|
||
while (inFlight > observed)
|
||
{
|
||
int previous = Interlocked.CompareExchange(ref _maxInFlight, inFlight, observed);
|
||
if (previous == observed)
|
||
{
|
||
return;
|
||
}
|
||
|
||
observed = previous;
|
||
}
|
||
}
|
||
}
|
||
}
|