fix(archreview): gateway core P0 remediation (GWC-01/02/03, TST-02, TST-12)
Interlocking changes across the gateway server (shared GatewaySession.cs / SessionManager.cs), committed together: - GWC-01 (Critical): alarm monitor now attaches as an internal (non-counted) distributor subscriber instead of a second raw drain of the single worker event channel; WorkerClient._events -> SingleReader with a claimed-once guard so a future dual-consumer regression throws loudly. - GWC-02 (High): faulted sessions are swept in CloseExpiredLeasesAsync (IsFaultedReapable + FaultedReason); new FaultedGraceSeconds (default 0). - GWC-03 (High): configurable MaxSparseArrayLength (default 1_000_000) enforced before allocation. - TST-02 (High, security): StreamEvents attach now enforces the opening key id -> PermissionDenied on owner mismatch. - TST-12 (Medium): CLAUDE.md retention-defaults sentence corrected. Verified: NonWindows build clean; targeted tests 135/135 on macOS, plus WorkerClientTests 18/18 on the Windows host.
This commit is contained in:
@@ -470,6 +470,20 @@ public sealed class AlarmFailoverEndToEndTests
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
[EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
|
||||
{
|
||||
if (workerEvent.Event is not null)
|
||||
{
|
||||
yield return workerEvent.Event;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool TryGetSession(string sessionId, [MaybeNullWhen(false)] out GatewaySession session)
|
||||
{
|
||||
|
||||
@@ -783,6 +783,20 @@ public sealed class GatewayAlarmMonitorProviderModeTests
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
[EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
|
||||
{
|
||||
if (workerEvent.Event is not null)
|
||||
{
|
||||
yield return workerEvent.Event;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool TryGetSession(string sessionId, [MaybeNullWhen(false)] out GatewaySession session)
|
||||
{
|
||||
|
||||
@@ -520,4 +520,32 @@ public sealed class GatewayOptionsValidatorTests
|
||||
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(null, options);
|
||||
Assert.True(result.Succeeded);
|
||||
}
|
||||
|
||||
/// <summary>Verifies a <see cref="EventOptions.MaxSparseArrayLength"/> below one fails validation.</summary>
|
||||
/// <param name="value">Sparse-array cap under test.</param>
|
||||
[Theory]
|
||||
[InlineData(0)]
|
||||
[InlineData(-1)]
|
||||
public void Validate_Fails_WhenMaxSparseArrayLengthBelowOne(int value)
|
||||
{
|
||||
GatewayOptions options = CloneWithEvents(
|
||||
ValidOptions(),
|
||||
new EventOptions { MaxSparseArrayLength = value });
|
||||
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(null, options);
|
||||
Assert.True(result.Failed);
|
||||
Assert.Contains(
|
||||
result.Failures!,
|
||||
f => f.Contains("MxGateway:Events:MaxSparseArrayLength"));
|
||||
}
|
||||
|
||||
/// <summary>Verifies a positive <see cref="EventOptions.MaxSparseArrayLength"/> within range passes validation.</summary>
|
||||
[Fact]
|
||||
public void Validate_Succeeds_WhenMaxSparseArrayLengthWithinRange()
|
||||
{
|
||||
GatewayOptions options = CloneWithEvents(
|
||||
ValidOptions(),
|
||||
new EventOptions { MaxSparseArrayLength = 1 });
|
||||
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(null, options);
|
||||
Assert.True(result.Succeeded);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -281,6 +281,14 @@ public sealed class DashboardSessionAdminServiceTests
|
||||
throw new NotSupportedException();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
throw new NotSupportedException();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
|
||||
@@ -38,6 +38,58 @@ public sealed class EventStreamServiceTests
|
||||
Assert.Equal(1, metrics.GetSnapshot().StreamDisconnects);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// TST-02 (owner-scoped attach): the API key that opened a session may attach its event
|
||||
/// stream — the caller key equals the session owner, so streaming proceeds normally.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task StreamEventsAsync_WhenCallerKeyMatchesOwner_Streams()
|
||||
{
|
||||
FakeWorkerClient workerClient = new();
|
||||
GatewaySession session = CreateReadySession(workerClient, ownerKeyId: "key-owner");
|
||||
EventStreamService service = CreateService(new FakeSessionManager(session));
|
||||
workerClient.Events.Add(CreateWorkerEvent(sequence: 5, MxEventFamily.OnDataChange));
|
||||
workerClient.CompleteAfterConfiguredEvents = true;
|
||||
|
||||
List<MxEvent> events = [];
|
||||
await foreach (MxEvent mxEvent in service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: "key-owner", CancellationToken.None)
|
||||
.WithCancellation(CancellationToken.None))
|
||||
{
|
||||
events.Add(mxEvent);
|
||||
}
|
||||
|
||||
Assert.Equal([5UL], events.Select(mxEvent => mxEvent.WorkerSequence).ToArray());
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// TST-02 (owner-scoped attach, security control): a caller whose API key differs from
|
||||
/// the key that opened the session is rejected with a <see cref="SessionManagerErrorCode.PermissionDenied"/>
|
||||
/// fault before any events are streamed — closing the reconnect/fan-out trust-boundary hole.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task StreamEventsAsync_WhenCallerKeyDiffersFromOwner_ThrowsPermissionDenied()
|
||||
{
|
||||
FakeWorkerClient workerClient = new();
|
||||
GatewaySession session = CreateReadySession(workerClient, ownerKeyId: "key-owner");
|
||||
EventStreamService service = CreateService(new FakeSessionManager(session));
|
||||
|
||||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(async () =>
|
||||
{
|
||||
await foreach (MxEvent _ in service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: "key-intruder", CancellationToken.None)
|
||||
.WithCancellation(CancellationToken.None))
|
||||
{
|
||||
// No event should be yielded — the owner check runs before the first attach.
|
||||
}
|
||||
});
|
||||
|
||||
Assert.Equal(SessionManagerErrorCode.PermissionDenied, exception.ErrorCode);
|
||||
Assert.Equal(0, session.ActiveEventSubscriberCount);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that a second event subscriber is rejected when one is already active.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
@@ -48,13 +100,13 @@ public sealed class EventStreamServiceTests
|
||||
EventStreamService service = CreateService(new FakeSessionManager(session));
|
||||
using CancellationTokenSource firstSubscriberCancellation = new();
|
||||
await using IAsyncEnumerator<MxEvent> firstSubscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), firstSubscriberCancellation.Token)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, firstSubscriberCancellation.Token)
|
||||
.GetAsyncEnumerator(firstSubscriberCancellation.Token);
|
||||
Task<bool> firstMoveTask = firstSubscriber.MoveNextAsync().AsTask();
|
||||
|
||||
await WaitUntilAsync(() => session.ActiveEventSubscriberCount == 1);
|
||||
await using IAsyncEnumerator<MxEvent> secondSubscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||||
@@ -78,7 +130,7 @@ public sealed class EventStreamServiceTests
|
||||
EventStreamService service = CreateService(new FakeSessionManager(session));
|
||||
using CancellationTokenSource cancellationTokenSource = new();
|
||||
await using IAsyncEnumerator<MxEvent> subscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), cancellationTokenSource.Token)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, cancellationTokenSource.Token)
|
||||
.GetAsyncEnumerator(cancellationTokenSource.Token);
|
||||
Task<bool> moveTask = subscriber.MoveNextAsync().AsTask();
|
||||
|
||||
@@ -108,7 +160,7 @@ public sealed class EventStreamServiceTests
|
||||
workerClient.Events.Add(CreateWorkerEvent(sequence: 3, MxEventFamily.OnDataChange));
|
||||
workerClient.CompleteAfterConfiguredEvents = true;
|
||||
await using IAsyncEnumerator<MxEvent> subscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
Assert.True(await subscriber.MoveNextAsync().AsTask().WaitAsync(TestTimeout));
|
||||
@@ -142,10 +194,10 @@ public sealed class EventStreamServiceTests
|
||||
firstWorkerClient.CompleteAfterConfiguredEvents = true;
|
||||
secondWorkerClient.CompleteAfterConfiguredEvents = true;
|
||||
await using IAsyncEnumerator<MxEvent> firstSubscriber = service
|
||||
.StreamEventsAsync(CreateRequest(firstSession.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(firstSession.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
await using IAsyncEnumerator<MxEvent> secondSubscriber = service
|
||||
.StreamEventsAsync(CreateRequest(secondSession.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(secondSession.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
Assert.True(await firstSubscriber.MoveNextAsync().AsTask().WaitAsync(TestTimeout));
|
||||
@@ -192,7 +244,7 @@ public sealed class EventStreamServiceTests
|
||||
|
||||
workerClient.CompleteAfterConfiguredEvents = true;
|
||||
await using IAsyncEnumerator<MxEvent> subscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
// The pump fans 50 events into a subscriber channel with capacity 1 faster than this
|
||||
@@ -246,7 +298,7 @@ public sealed class EventStreamServiceTests
|
||||
|
||||
workerClient.CompleteAfterConfiguredEvents = true;
|
||||
await using IAsyncEnumerator<MxEvent> subscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
SessionManagerException exception = await Assert.ThrowsAsync<SessionManagerException>(
|
||||
@@ -298,7 +350,7 @@ public sealed class EventStreamServiceTests
|
||||
using GatewayMetrics metrics = new();
|
||||
EventStreamService service = CreateService(new FakeSessionManager(session), metrics);
|
||||
await using IAsyncEnumerator<MxEvent> subscriber = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
WorkerClientException exception = await Assert.ThrowsAsync<WorkerClientException>(
|
||||
@@ -333,7 +385,7 @@ public sealed class EventStreamServiceTests
|
||||
|
||||
// Resume after sequence 2: retained window [1..5] covers it — replay 3,4,5 then live.
|
||||
await using IAsyncEnumerator<MxEvent> resume = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 2), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 2), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
MxEvent r3 = await ReadNextAsync(resume);
|
||||
@@ -374,7 +426,7 @@ public sealed class EventStreamServiceTests
|
||||
// Resume after 1: events 1,2 are below the oldest retained (3) and were evicted, so
|
||||
// they are unrecoverable => sentinel first, then the retained tail 3,4,5, then live.
|
||||
await using IAsyncEnumerator<MxEvent> realResume = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 1), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 1), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
MxEvent sentinel = await ReadNextAsync(realResume);
|
||||
@@ -419,7 +471,7 @@ public sealed class EventStreamServiceTests
|
||||
// Resume after 2: replay 3,4 then live 5,6,7. Collect across the boundary and assert
|
||||
// the full sequence is contiguous with no duplicate and no skip.
|
||||
await using IAsyncEnumerator<MxEvent> resume = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 2), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 2), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
List<ulong> collected = [];
|
||||
@@ -461,7 +513,7 @@ public sealed class EventStreamServiceTests
|
||||
// reads must be exactly 4 then 5 (no sentinel, no <=3 event); a live tag confirms the
|
||||
// stream resumed live strictly after 5.
|
||||
await using IAsyncEnumerator<MxEvent> resume = service
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 3), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(session.SessionId, afterWorkerSequence: 3), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
|
||||
MxEvent first = await ReadNextAsync(resume);
|
||||
@@ -512,7 +564,7 @@ public sealed class EventStreamServiceTests
|
||||
int expectedCount)
|
||||
{
|
||||
await using IAsyncEnumerator<MxEvent> primer = service
|
||||
.StreamEventsAsync(CreateRequest(sessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(sessionId), callerKeyId: null, CancellationToken.None)
|
||||
.GetAsyncEnumerator();
|
||||
for (int i = 0; i < expectedCount; i++)
|
||||
{
|
||||
@@ -551,7 +603,7 @@ public sealed class EventStreamServiceTests
|
||||
{
|
||||
List<MxEvent> events = [];
|
||||
await foreach (MxEvent mxEvent in service
|
||||
.StreamEventsAsync(CreateRequest(sessionId), CancellationToken.None)
|
||||
.StreamEventsAsync(CreateRequest(sessionId), callerKeyId: null, CancellationToken.None)
|
||||
.WithCancellation(CancellationToken.None))
|
||||
{
|
||||
events.Add(mxEvent);
|
||||
@@ -575,7 +627,8 @@ public sealed class EventStreamServiceTests
|
||||
int queueCapacity = 8,
|
||||
GatewayMetrics? metrics = null,
|
||||
EventBackpressurePolicy backpressurePolicy = EventBackpressurePolicy.FailFast,
|
||||
int replayBufferCapacity = 1024)
|
||||
int replayBufferCapacity = 1024,
|
||||
string? ownerKeyId = null)
|
||||
{
|
||||
// The per-subscriber overflow policy now lives in the session's
|
||||
// SessionEventDistributor, so the session must share the same metrics sink and
|
||||
@@ -587,7 +640,7 @@ public sealed class EventStreamServiceTests
|
||||
"pipe",
|
||||
"nonce",
|
||||
"client",
|
||||
ownerKeyId: null,
|
||||
ownerKeyId: ownerKeyId,
|
||||
"client-session",
|
||||
"client-correlation",
|
||||
TimeSpan.FromSeconds(30),
|
||||
@@ -702,6 +755,14 @@ public sealed class EventStreamServiceTests
|
||||
return _sessions[sessionId].ReadEventsAsync(cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
throw new NotSupportedException();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
|
||||
@@ -948,6 +948,14 @@ public sealed class MxAccessGatewayServiceConstraintTests
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
throw new NotSupportedException();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
@@ -994,6 +1002,7 @@ public sealed class MxAccessGatewayServiceConstraintTests
|
||||
/// <inheritdoc />
|
||||
public async IAsyncEnumerable<MxEvent> StreamEventsAsync(
|
||||
StreamEventsRequest request,
|
||||
string? callerKeyId,
|
||||
[System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (WorkerEvent ev in sessionManager.Events)
|
||||
|
||||
@@ -616,6 +616,14 @@ public sealed class MxAccessGatewayServiceTests
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
throw new NotSupportedException();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
@@ -653,6 +661,7 @@ public sealed class MxAccessGatewayServiceTests
|
||||
/// <inheritdoc />
|
||||
public async IAsyncEnumerable<MxEvent> StreamEventsAsync(
|
||||
StreamEventsRequest request,
|
||||
string? callerKeyId,
|
||||
[EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
sessionManager.RecordReadEventsSessionId(request.SessionId);
|
||||
|
||||
+78
-1
@@ -88,7 +88,7 @@ public sealed class GatewaySessionDashboardMirrorTests
|
||||
Task grpcReader = Task.Run(async () =>
|
||||
{
|
||||
await foreach (MxEvent mxEvent in service
|
||||
.StreamEventsAsync(new StreamEventsRequest { SessionId = session.SessionId }, CancellationToken.None)
|
||||
.StreamEventsAsync(new StreamEventsRequest { SessionId = session.SessionId }, callerKeyId: null, CancellationToken.None)
|
||||
.WithCancellation(CancellationToken.None))
|
||||
{
|
||||
grpcEvents.Add(mxEvent);
|
||||
@@ -107,6 +107,65 @@ public sealed class GatewaySessionDashboardMirrorTests
|
||||
Assert.Equal([1UL, 2UL, 3UL], broadcaster.Captures.Select(capture => capture.MxEvent.WorkerSequence).ToArray());
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// GWC-01 regression: with the internal dashboard mirror active, a second internal
|
||||
/// subscriber (the alarm monitor's feed, attached via
|
||||
/// <see cref="GatewaySession.AttachInternalEventSubscriber"/>) receives EVERY event —
|
||||
/// including the alarm <c>Acknowledge</c> transition — rather than the two consumers
|
||||
/// each getting a random half of the single worker channel. Before the fix the alarm
|
||||
/// monitor drained the worker channel directly, so with the dashboard mirror pump also
|
||||
/// draining it the two split the stream and Acknowledge transitions were silently lost.
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task InternalAlarmSubscriber_AndDashboardMirror_BothReceiveEveryEvent()
|
||||
{
|
||||
FakeWorkerClient workerClient = new();
|
||||
workerClient.Events.Add(CreateWorkerEvent(1, MxEventFamily.OnDataChange));
|
||||
workerClient.Events.Add(CreateAlarmTransitionEvent(2, AlarmTransitionKind.Raise));
|
||||
workerClient.Events.Add(CreateAlarmTransitionEvent(3, AlarmTransitionKind.Acknowledge));
|
||||
workerClient.Events.Add(CreateAlarmTransitionEvent(4, AlarmTransitionKind.Clear));
|
||||
workerClient.CompleteAfterConfiguredEvents = true;
|
||||
// Hold the finite stream until BOTH internal subscribers (dashboard mirror + alarm feed)
|
||||
// are registered so neither misses an event before the pump drains.
|
||||
workerClient.HoldEventsUntilReleased();
|
||||
RecordingDashboardEventBroadcaster broadcaster = new();
|
||||
|
||||
await using GatewaySession session = CreateSession(workerClient, broadcaster);
|
||||
session.AttachWorkerClient(workerClient);
|
||||
|
||||
// MarkReady registers the internal dashboard subscriber and starts the (gated) pump.
|
||||
session.MarkReady();
|
||||
|
||||
// Attach the alarm monitor's internal subscriber and drain it concurrently. Registered
|
||||
// BEFORE the stream is released, so it is present at pump start alongside the dashboard.
|
||||
using IEventSubscriberLease alarmLease = session.AttachInternalEventSubscriber();
|
||||
List<MxEvent> alarmEvents = [];
|
||||
Task alarmReader = Task.Run(async () =>
|
||||
{
|
||||
await foreach (MxEvent mxEvent in alarmLease.Reader.ReadAllAsync(CancellationToken.None))
|
||||
{
|
||||
alarmEvents.Add(mxEvent);
|
||||
}
|
||||
});
|
||||
|
||||
workerClient.ReleaseEvents();
|
||||
|
||||
await WaitUntilAsync(() => broadcaster.Captures.Count == 4 && alarmEvents.Count == 4);
|
||||
await alarmReader.WaitAsync(TestTimeout);
|
||||
|
||||
// Both internal consumers see the full, identical, ordered stream — no splitting.
|
||||
Assert.Equal([1UL, 2UL, 3UL, 4UL], broadcaster.Captures.Select(capture => capture.MxEvent.WorkerSequence).ToArray());
|
||||
Assert.Equal([1UL, 2UL, 3UL, 4UL], alarmEvents.Select(mxEvent => mxEvent.WorkerSequence).ToArray());
|
||||
|
||||
// The Acknowledge transition specifically reaches the alarm feed (the transition the
|
||||
// pre-fix split silently dropped, leaving clients showing unacked alarms indefinitely).
|
||||
Assert.Contains(
|
||||
alarmEvents,
|
||||
mxEvent => mxEvent.BodyCase == MxEvent.BodyOneofCase.OnAlarmTransition
|
||||
&& mxEvent.OnAlarmTransition.TransitionKind == AlarmTransitionKind.Acknowledge);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Hazard guard: starting the pump at Ready with a fast-completing worker stream
|
||||
/// and zero subscribers used to drain into nothing and leave a later subscriber hanging.
|
||||
@@ -222,6 +281,19 @@ public sealed class GatewaySessionDashboardMirrorTests
|
||||
return new WorkerEvent { Event = mxEvent };
|
||||
}
|
||||
|
||||
private static WorkerEvent CreateAlarmTransitionEvent(ulong sequence, AlarmTransitionKind kind)
|
||||
{
|
||||
MxEvent mxEvent = new()
|
||||
{
|
||||
SessionId = "session-dashboard-mirror",
|
||||
Family = MxEventFamily.OnAlarmTransition,
|
||||
WorkerSequence = sequence,
|
||||
OnAlarmTransition = new OnAlarmTransitionEvent { TransitionKind = kind },
|
||||
};
|
||||
|
||||
return new WorkerEvent { Event = mxEvent };
|
||||
}
|
||||
|
||||
private static async Task WaitUntilAsync(Func<bool> predicate, [CallerArgumentExpression(nameof(predicate))] string? condition = null)
|
||||
{
|
||||
using CancellationTokenSource cancellationTokenSource = new(TestTimeout);
|
||||
@@ -280,6 +352,11 @@ public sealed class GatewaySessionDashboardMirrorTests
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken) => session.ReadEventsAsync(cancellationToken);
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken) => throw new NotSupportedException();
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
|
||||
@@ -975,6 +975,68 @@ public sealed class SessionManagerTests
|
||||
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>Verifies that shutdown closes all registered sessions.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
@@ -1023,7 +1085,8 @@ public sealed class SessionManagerTests
|
||||
int maxSessions = 64,
|
||||
int defaultLeaseSeconds = 1800,
|
||||
int detachGraceSeconds = 0,
|
||||
int workerReadyWaitTimeoutMs = 0)
|
||||
int workerReadyWaitTimeoutMs = 0,
|
||||
int faultedGraceSeconds = 0)
|
||||
{
|
||||
return new GatewayOptions
|
||||
{
|
||||
@@ -1034,6 +1097,7 @@ public sealed class SessionManagerTests
|
||||
DefaultLeaseSeconds = defaultLeaseSeconds,
|
||||
DetachGraceSeconds = detachGraceSeconds,
|
||||
WorkerReadyWaitTimeoutMs = workerReadyWaitTimeoutMs,
|
||||
FaultedGraceSeconds = faultedGraceSeconds,
|
||||
},
|
||||
Worker = new WorkerOptions
|
||||
{
|
||||
|
||||
@@ -222,4 +222,35 @@ public sealed class SparseArrayExpanderTests
|
||||
RpcException ex = Assert.Throws<RpcException>(() => SparseArrayExpander.Expand(value));
|
||||
Assert.Equal(StatusCode.InvalidArgument, ex.StatusCode);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that a total length above the configured cap throws <see cref="StatusCode.InvalidArgument"/> before the full array is allocated.</summary>
|
||||
[Fact]
|
||||
public void Expand_TotalLengthExceedsConfiguredCap_ThrowsBeforeAllocation()
|
||||
{
|
||||
// A total_length that would force a huge allocation, but well below Array.MaxLength,
|
||||
// so only the configured cap can reject it (the Array.MaxLength backstop would not).
|
||||
MxValue value = SparseValue(MxDataType.Integer, 500_000_000u);
|
||||
|
||||
RpcException ex = Assert.Throws<RpcException>(() => SparseArrayExpander.Expand(value, maxSparseArrayLength: 1_000_000));
|
||||
Assert.Equal(StatusCode.InvalidArgument, ex.StatusCode);
|
||||
Assert.Contains("MaxSparseArrayLength", ex.Status.Detail, StringComparison.Ordinal);
|
||||
|
||||
// The value must be untouched — expansion (allocation) never ran.
|
||||
Assert.Equal(MxValue.KindOneofCase.SparseArrayValue, value.KindCase);
|
||||
}
|
||||
|
||||
/// <summary>Verifies that a total length at or below the configured cap still expands normally.</summary>
|
||||
[Fact]
|
||||
public void Expand_TotalLengthAtConfiguredCap_Expands()
|
||||
{
|
||||
MxValue value = SparseValue(
|
||||
MxDataType.Integer,
|
||||
4,
|
||||
(1, new MxValue { Int32Value = 7 }));
|
||||
|
||||
SparseArrayExpander.Expand(value, maxSparseArrayLength: 4);
|
||||
|
||||
Assert.Equal(MxValue.KindOneofCase.ArrayValue, value.KindCase);
|
||||
Assert.Equal(new[] { 0, 7, 0, 0 }, value.ArrayValue.Int32Values.Values);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -122,6 +122,27 @@ public sealed class WorkerClientTests
|
||||
Assert.Equal(MxEventFamily.OperationComplete, events.Current.Event.Family);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The worker event channel is single-reader: a second <see cref="WorkerClient.ReadEventsAsync"/>
|
||||
/// enumerator must throw rather than silently split events between two consumers (GWC-01).
|
||||
/// </summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
public async Task ReadEventsAsync_SecondEnumerator_Throws()
|
||||
{
|
||||
await using PipePair pipePair = await PipePair.CreateAsync();
|
||||
await using WorkerClient client = CreateClient(pipePair);
|
||||
await CompleteHandshakeAsync(client, pipePair);
|
||||
using CancellationTokenSource cancellationTokenSource = new(TestTimeout);
|
||||
|
||||
// First call claims the single reader.
|
||||
IAsyncEnumerable<WorkerEvent> first = client.ReadEventsAsync(cancellationTokenSource.Token);
|
||||
await using IAsyncEnumerator<WorkerEvent> firstEnumerator = first.GetAsyncEnumerator(cancellationTokenSource.Token);
|
||||
|
||||
// A second call must fail loudly at call time rather than racing the first for events.
|
||||
Assert.Throws<InvalidOperationException>(() => client.ReadEventsAsync(cancellationTokenSource.Token));
|
||||
}
|
||||
|
||||
/// <summary>Verifies that the read loop faults the client when the event queue overflows.</summary>
|
||||
/// <returns>A task that represents the asynchronous operation.</returns>
|
||||
[Fact]
|
||||
|
||||
+9
@@ -477,6 +477,14 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
|
||||
return AsyncEnumerable.Empty<WorkerEvent>();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
|
||||
string sessionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return AsyncEnumerable.Empty<MxEvent>();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<SessionCloseResult> CloseSessionAsync(
|
||||
string sessionId,
|
||||
@@ -515,6 +523,7 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
|
||||
/// <inheritdoc />
|
||||
public async IAsyncEnumerable<MxEvent> StreamEventsAsync(
|
||||
StreamEventsRequest request,
|
||||
string? callerKeyId,
|
||||
[EnumeratorCancellation] CancellationToken cancellationToken)
|
||||
{
|
||||
await Task.CompletedTask;
|
||||
|
||||
Reference in New Issue
Block a user