e6842c108a
Root-caused from a full-suite-load-only failure of GrpcStreamIntegrationTests.Pipeline_DuplicateCorrelationId_ReplacesStream: ObjectDisposedException escaping SubscribeInstance. SubscribeInstance's duplicate-prevention path cancelled AND disposed the replaced stream's CancellationTokenSource. That CTS belongs to the replaced handler's own `using var streamCts`, which is still running and still has to read `streamCts.Token` — the Dispose raced that read. The race is pre-existing and independent of R2: the same Dispose sat against the same first-token-read when the handler still used `ReadAllAsync(streamCts.Token)`; it only ever loses under enough scheduling pressure to land the replacement inside the first stream's setup window. Cancel only. The owning handler's `using` still disposes it exactly once on every exit path, and cancellation is all replacement ever needed. The Cancel is wrapped for the converse race (owner already finished and disposed), mirroring CancelAllStreams(). Regression test makes the race deterministic by gating the first stream inside its setup window — the _activeStreams entry is registered before Subscribe is called, so the replacement always lands before the first stream reads its token. Verified fail-before (ObjectDisposedException) / pass-after by reinstating the Dispose as a negative control.
857 lines
34 KiB
C#
857 lines
34 KiB
C#
using System.Diagnostics.Metrics;
|
|
using System.Threading.Channels;
|
|
using Akka.Actor;
|
|
using Akka.TestKit.Xunit2;
|
|
using Grpc.Core;
|
|
using Microsoft.Extensions.Logging;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
using NSubstitute;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Observability;
|
|
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
|
|
|
|
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests.Grpc;
|
|
|
|
public class SiteStreamGrpcServerTests : TestKit
|
|
{
|
|
private readonly ISiteStreamSubscriber _subscriber;
|
|
private readonly ILogger<SiteStreamGrpcServer> _logger;
|
|
|
|
public SiteStreamGrpcServerTests()
|
|
{
|
|
_subscriber = Substitute.For<ISiteStreamSubscriber>();
|
|
_subscriber.Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>())
|
|
.Returns("sub-1");
|
|
_subscriber.SubscribeSiteAlarms(Arg.Any<IActorRef>())
|
|
.Returns("site-sub-1");
|
|
_logger = NullLogger<SiteStreamGrpcServer>.Instance;
|
|
}
|
|
|
|
private SiteStreamGrpcServer CreateServer(int maxStreams = 100)
|
|
{
|
|
return new SiteStreamGrpcServer(_subscriber, _logger, maxStreams);
|
|
}
|
|
|
|
private static InstanceStreamRequest MakeRequest(string correlationId = "corr-1", string instance = "Site1.Pump01")
|
|
{
|
|
return new InstanceStreamRequest
|
|
{
|
|
CorrelationId = correlationId,
|
|
InstanceUniqueName = instance
|
|
};
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RejectsWhenNotReady()
|
|
{
|
|
var server = CreateServer();
|
|
// Do NOT call SetReady()
|
|
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context = CreateMockContext();
|
|
|
|
var ex = await Assert.ThrowsAsync<RpcException>(
|
|
() => server.SubscribeInstance(MakeRequest(), writer, context));
|
|
|
|
Assert.Equal(StatusCode.Unavailable, ex.StatusCode);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RejectsWhenMaxStreamsReached()
|
|
{
|
|
var server = CreateServer(maxStreams: 1);
|
|
server.SetReady(Sys);
|
|
|
|
// Start one stream that blocks
|
|
var cts1 = new CancellationTokenSource();
|
|
var context1 = CreateMockContext(cts1.Token);
|
|
var writer1 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var stream1Task = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-1"), writer1, context1));
|
|
|
|
// Wait for the first stream to register
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Second stream should be rejected
|
|
var writer2 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context2 = CreateMockContext();
|
|
|
|
var ex = await Assert.ThrowsAsync<RpcException>(
|
|
() => server.SubscribeInstance(MakeRequest("corr-2"), writer2, context2));
|
|
|
|
Assert.Equal(StatusCode.ResourceExhausted, ex.StatusCode);
|
|
|
|
// Clean up first stream
|
|
cts1.Cancel();
|
|
await stream1Task;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task CancelsDuplicateCorrelationId()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts1 = new CancellationTokenSource();
|
|
var context1 = CreateMockContext(cts1.Token);
|
|
var writer1 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
// Start first stream
|
|
var stream1Task = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-dup"), writer1, context1));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Start second stream with same correlationId -- should cancel first
|
|
var cts2 = new CancellationTokenSource();
|
|
var context2 = CreateMockContext(cts2.Token);
|
|
var writer2 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var stream2Task = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-dup"), writer2, context2));
|
|
|
|
// First stream should complete (cancelled by duplicate replacement)
|
|
await stream1Task;
|
|
|
|
// Second stream should be active
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Clean up
|
|
cts2.Cancel();
|
|
await stream2Task;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task CleansUpOnCancellation()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-cleanup"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
// --- Host-017 / REQ-HOST-7: site-shutdown ordering ---
|
|
|
|
[Fact]
|
|
public async Task Host017_CancelAllStreams_CancelsActiveStreamsAndRefusesNewOnes()
|
|
{
|
|
// REQ-HOST-7 step (1)+(2): on CoordinatedShutdown the gRPC server must
|
|
// stop accepting new streams AND cancel every active stream so the
|
|
// client observes a clean Cancelled (not a silent stream that only
|
|
// times out via keepalive). Program.cs registers
|
|
// ApplicationStopping → CancelAllStreams(); this test exercises the
|
|
// server-side guarantee in isolation.
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts1 = new CancellationTokenSource();
|
|
var context1 = CreateMockContext(cts1.Token);
|
|
var writer1 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var stream1Task = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-shutdown-1"), writer1, context1));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Begin shutdown — flip the flag AND cancel the active stream.
|
|
server.CancelAllStreams();
|
|
|
|
Assert.True(server.IsShuttingDown);
|
|
|
|
// Active stream's await foreach observes OCE and falls through finally
|
|
// → entry is removed from _activeStreams.
|
|
await stream1Task;
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
|
|
// A second SubscribeInstance after shutdown is refused immediately
|
|
// with Unavailable rather than allowed to register a new stream.
|
|
var writer2 = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context2 = CreateMockContext();
|
|
var ex = await Assert.ThrowsAsync<RpcException>(
|
|
() => server.SubscribeInstance(MakeRequest("corr-shutdown-2"), writer2, context2));
|
|
Assert.Equal(StatusCode.Unavailable, ex.StatusCode);
|
|
Assert.Contains("shutting", ex.Status.Detail, StringComparison.OrdinalIgnoreCase);
|
|
}
|
|
|
|
[Fact]
|
|
public void Host017_CancelAllStreams_IsIdempotent()
|
|
{
|
|
// Repeated calls during a double-fire shutdown sequence must not throw.
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
server.CancelAllStreams();
|
|
server.CancelAllStreams();
|
|
|
|
Assert.True(server.IsShuttingDown);
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SubscribesAndRemovesFromStreamManager()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-sub", "Site1.Motor01"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Verify Subscribe was called
|
|
_subscriber.Received(1).Subscribe("Site1.Motor01", Arg.Any<IActorRef>());
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
// Verify RemoveSubscriber was called
|
|
_subscriber.Received(1).RemoveSubscriber(Arg.Any<IActorRef>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task WritesEventsToResponseStream()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
// Capture the relay actor so we can send it events
|
|
IActorRef? capturedActor = null;
|
|
_subscriber.Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>())
|
|
.Returns(ci =>
|
|
{
|
|
capturedActor = ci.Arg<IActorRef>();
|
|
return "sub-write";
|
|
});
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var writtenEvents = new List<SiteStreamEvent>();
|
|
writer.WriteAsync(Arg.Any<SiteStreamEvent>(), Arg.Any<CancellationToken>())
|
|
.Returns(Task.CompletedTask)
|
|
.AndDoes(ci => writtenEvents.Add(ci.Arg<SiteStreamEvent>()));
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-write", "Site1.Pump01"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => capturedActor != null);
|
|
|
|
// Send a domain event to the relay actor
|
|
var ts = DateTimeOffset.UtcNow;
|
|
capturedActor!.Tell(new Commons.Messages.Streaming.AttributeValueChanged(
|
|
"Site1.Pump01", "Path", "Attr", 99.5, "Good", ts));
|
|
|
|
// Wait for event to be written
|
|
await WaitForConditionAsync(() => writtenEvents.Count >= 1);
|
|
|
|
Assert.Single(writtenEvents);
|
|
Assert.Equal("corr-write", writtenEvents[0].CorrelationId);
|
|
Assert.Equal(SiteStreamEvent.EventOneofCase.AttributeChanged, writtenEvents[0].EventCase);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
}
|
|
|
|
// --- SubscribeSite (site-wide, alarm-only aggregated stream, plan #10 T2) ---
|
|
|
|
private static SiteStreamRequest MakeSiteRequest(string correlationId = "site-corr-1")
|
|
=> new() { CorrelationId = correlationId };
|
|
|
|
[Fact]
|
|
public async Task SubscribeSite_SubscribesSiteAlarmsAndRemovesOnCancel()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeSite(
|
|
MakeSiteRequest("site-corr-sub"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// Site-wide handler must call SubscribeSiteAlarms (no instance filter),
|
|
// never the per-instance Subscribe.
|
|
_subscriber.Received(1).SubscribeSiteAlarms(Arg.Any<IActorRef>());
|
|
_subscriber.DidNotReceive().Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>());
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
_subscriber.Received(1).RemoveSubscriber(Arg.Any<IActorRef>());
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SubscribeSite_RejectsUnsafeCorrelationId()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context = CreateMockContext();
|
|
|
|
var ex = await Assert.ThrowsAsync<RpcException>(
|
|
() => server.SubscribeSite(MakeSiteRequest("bad/id"), writer, context));
|
|
|
|
Assert.Equal(StatusCode.InvalidArgument, ex.StatusCode);
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SubscribeSite_RelaysAlarmStateChangedAsAlarmStateUpdate()
|
|
{
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
// Capture the relay actor spawned for the site-wide subscription.
|
|
IActorRef? capturedActor = null;
|
|
_subscriber.SubscribeSiteAlarms(Arg.Any<IActorRef>())
|
|
.Returns(ci =>
|
|
{
|
|
capturedActor = ci.Arg<IActorRef>();
|
|
return "site-sub-relay";
|
|
});
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var writtenEvents = new List<SiteStreamEvent>();
|
|
writer.WriteAsync(Arg.Any<SiteStreamEvent>(), Arg.Any<CancellationToken>())
|
|
.Returns(Task.CompletedTask)
|
|
.AndDoes(ci => writtenEvents.Add(ci.Arg<SiteStreamEvent>()));
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeSite(
|
|
MakeSiteRequest("site-corr-write"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => capturedActor != null);
|
|
|
|
// A real alarm transition for ANY instance must arrive as an AlarmStateUpdate.
|
|
capturedActor!.Tell(new Commons.Messages.Streaming.AlarmStateChanged(
|
|
"Site1.Pump01",
|
|
"HighPressure",
|
|
Commons.Types.Enums.AlarmState.Active,
|
|
700,
|
|
DateTimeOffset.UtcNow));
|
|
|
|
await WaitForConditionAsync(() => writtenEvents.Count >= 1);
|
|
|
|
Assert.Single(writtenEvents);
|
|
Assert.Equal("site-corr-write", writtenEvents[0].CorrelationId);
|
|
Assert.Equal(SiteStreamEvent.EventOneofCase.AlarmChanged, writtenEvents[0].EventCase);
|
|
Assert.Equal("Site1.Pump01", writtenEvents[0].AlarmChanged.InstanceUniqueName);
|
|
Assert.Equal("HighPressure", writtenEvents[0].AlarmChanged.AlarmName);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
}
|
|
|
|
[Theory]
|
|
[InlineData("corr/with/slash")]
|
|
[InlineData("corr with space")]
|
|
[InlineData("")]
|
|
[InlineData("$weird")]
|
|
public async Task RejectsCorrelationIdThatIsNotActorNameSafe(string badCorrelationId)
|
|
{
|
|
// Communication-014 regression: a public gRPC SubscribeInstance must not feed
|
|
// an untrusted correlation_id straight into an Akka actor name. An unsafe id
|
|
// must be rejected cleanly with InvalidArgument rather than escaping as an
|
|
// unhandled InvalidActorNameException.
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context = CreateMockContext();
|
|
|
|
var ex = await Assert.ThrowsAsync<RpcException>(
|
|
() => server.SubscribeInstance(MakeRequest(badCorrelationId), writer, context));
|
|
|
|
Assert.Equal(StatusCode.InvalidArgument, ex.StatusCode);
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task AcceptsActorNameSafeCorrelationId()
|
|
{
|
|
// A normal GUID-style correlation id (what central always supplies) is accepted.
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest(Guid.NewGuid().ToString()), writer, context));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Comm021_SubscribeThrows_StopsRelayActorAndRemovesActiveStreamEntry()
|
|
{
|
|
// Communication-021 regression: SubscribeInstance creates a StreamRelayActor
|
|
// and registers an _activeStreams entry BEFORE calling _streamSubscriber.Subscribe.
|
|
// If Subscribe throws (e.g. stale instance, site runtime shutting down) and the
|
|
// pre-fix code lets the throw escape without the wrapping try, the relay actor
|
|
// and the activeStreams entry both leak. The fix wraps the Subscribe call so the
|
|
// catch deterministically stops the actor and removes the entry before re-throw.
|
|
var subscriber = Substitute.For<ISiteStreamSubscriber>();
|
|
subscriber.Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>())
|
|
.Returns<string>(_ => throw new InvalidOperationException("instance not found"));
|
|
|
|
var server = new SiteStreamGrpcServer(subscriber, _logger);
|
|
server.SetReady(Sys);
|
|
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
var context = CreateMockContext();
|
|
|
|
// The InvalidOperationException is expected to propagate (the gRPC stack maps
|
|
// unhandled throws to Internal); the load-bearing assertion is the cleanup.
|
|
await Assert.ThrowsAsync<InvalidOperationException>(
|
|
() => server.SubscribeInstance(MakeRequest("corr-comm021"), writer, context));
|
|
|
|
// _activeStreams entry was inserted before Subscribe was called; the catch
|
|
// must remove it so a follow-up subscription with the same correlation id is
|
|
// not blocked, and the relay actor must be stopped so it does not leak.
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
|
|
// RemoveSubscriber must NOT have been called (Subscribe never returned a
|
|
// subscription id) — verifying we hit the catch path, not the finally path.
|
|
subscriber.DidNotReceive().RemoveSubscriber(Arg.Any<IActorRef>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SiteConnectionUpGauge_GoesToOneOnConnect_AndBackToZeroOnCancel()
|
|
{
|
|
// Telemetry follow-on: the scadabridge.site.connection.up gauge must read
|
|
// exactly 1 while a site stream is established and return to 0 once the
|
|
// stream terminates on the cancel path — proving SiteConnectionOpened() is
|
|
// matched by exactly one SiteConnectionClosed() in the handler's finally.
|
|
var server = CreateServer();
|
|
server.SetReady(Sys);
|
|
|
|
long ReadGauge()
|
|
{
|
|
long observed = 0;
|
|
using var listener = new MeterListener();
|
|
listener.InstrumentPublished = (instrument, l) =>
|
|
{
|
|
if (instrument.Meter.Name == ScadaBridgeTelemetry.MeterName &&
|
|
instrument.Name == "scadabridge.site.connection.up")
|
|
{
|
|
l.EnableMeasurementEvents(instrument);
|
|
}
|
|
};
|
|
listener.SetMeasurementEventCallback<long>((_, measurement, _, _) => observed = measurement);
|
|
listener.Start();
|
|
listener.RecordObservableInstruments();
|
|
return observed;
|
|
}
|
|
|
|
var baseline = ReadGauge();
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-gauge", "Site1.Pump01"), writer, context));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
|
|
// While the stream is up the gauge is one above whatever baseline other
|
|
// (possibly parallel) tests left behind — read relative so the assertion
|
|
// is robust to test interleaving on the process-wide static counter.
|
|
//
|
|
// SiteConnectionOpened() runs AFTER the _activeStreams insertion that
|
|
// WaitForConditionAsync(ActiveStreamCount == 1) keys off (see
|
|
// SiteStreamGrpcServer.SubscribeInstance), so a one-shot ReadGauge() here
|
|
// can race the increment under full-suite CPU oversubscription. Poll the
|
|
// gauge with a generous timeout until it reaches baseline + 1.
|
|
AwaitAssert(() => Assert.Equal(baseline + 1, ReadGauge()), TimeSpan.FromSeconds(5));
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 0);
|
|
|
|
// After the cancel path runs the finally, the gauge is balanced back to
|
|
// the baseline — no leaked "up" count. Poll here too: the gauge decrement
|
|
// and the stream-count drop are independent observations, so under load
|
|
// give the finally's SiteConnectionClosed() time to land.
|
|
AwaitAssert(() => Assert.Equal(baseline, ReadGauge()), TimeSpan.FromSeconds(5));
|
|
}
|
|
|
|
[Fact]
|
|
public void SetReady_AllowsStreamCreation()
|
|
{
|
|
var server = CreateServer();
|
|
// Initially not ready -- just verify the property works
|
|
server.SetReady(Sys);
|
|
// No assertion needed -- the other tests verify that SetReady enables streaming
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
private static ServerCallContext CreateMockContext(CancellationToken cancellationToken = default)
|
|
{
|
|
var context = Substitute.For<ServerCallContext>();
|
|
context.CancellationToken.Returns(cancellationToken);
|
|
return context;
|
|
}
|
|
|
|
private static async Task WaitForConditionAsync(Func<bool> condition, int timeoutMs = 5000)
|
|
{
|
|
var deadline = DateTime.UtcNow.AddMilliseconds(timeoutMs);
|
|
while (!condition() && DateTime.UtcNow < deadline)
|
|
{
|
|
await Task.Delay(25);
|
|
}
|
|
|
|
Assert.True(condition(), $"Condition not met within {timeoutMs}ms");
|
|
}
|
|
|
|
// ── WP2.3: the site-wide alarm feed gets its OWN, larger send channel ──
|
|
|
|
[Fact]
|
|
public void SiteAlarmStream_HasItsOwnLargerChannel_ThanTheDebugView()
|
|
{
|
|
// Sharing the Debug View's 1000-slot DropOldest channel meant an alarm burst during
|
|
// a WAN stall silently evicted operator-visible transitions to make room for
|
|
// diagnostics traffic. The two feeds are now sized independently.
|
|
var options = Microsoft.Extensions.Options.Options.Create(new CommunicationOptions());
|
|
var server = new SiteStreamGrpcServer(_subscriber, _logger, options);
|
|
|
|
Assert.Equal(1000, server.InstanceChannelCapacity);
|
|
Assert.Equal(20_000, server.SiteAlarmChannelCapacity);
|
|
Assert.True(server.SiteAlarmChannelCapacity > server.InstanceChannelCapacity);
|
|
}
|
|
|
|
[Fact]
|
|
public void ChannelCapacities_AreBoundFromOptions_AndFloorAtOne()
|
|
{
|
|
var options = Microsoft.Extensions.Options.Options.Create(new CommunicationOptions
|
|
{
|
|
GrpcInstanceStreamChannelCapacity = 42,
|
|
GrpcSiteAlarmStreamChannelCapacity = 4242,
|
|
});
|
|
var server = new SiteStreamGrpcServer(_subscriber, _logger, options);
|
|
|
|
Assert.Equal(42, server.InstanceChannelCapacity);
|
|
Assert.Equal(4242, server.SiteAlarmChannelCapacity);
|
|
|
|
// A misconfigured zero/negative capacity must not throw at channel-construction
|
|
// time deep inside a live RPC — it floors at one instead.
|
|
var degenerate = new SiteStreamGrpcServer(_subscriber, _logger,
|
|
Microsoft.Extensions.Options.Options.Create(new CommunicationOptions
|
|
{
|
|
GrpcInstanceStreamChannelCapacity = 0,
|
|
GrpcSiteAlarmStreamChannelCapacity = -5,
|
|
}));
|
|
Assert.Equal(1, degenerate.InstanceChannelCapacity);
|
|
Assert.Equal(1, degenerate.SiteAlarmChannelCapacity);
|
|
}
|
|
|
|
[Fact]
|
|
public void DroppedStreamEventCount_StartsAtZero()
|
|
{
|
|
// The raw counter behind scadabridge.site.stream.events_dropped — a fresh node has
|
|
// evicted nothing.
|
|
var server = CreateServer();
|
|
Assert.Equal(0, server.DroppedStreamEventCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task DuplicateReplacement_CancelsTheReplacedStream_WithoutDisposingItsCts()
|
|
{
|
|
// Regression: the duplicate-replacement path used to Cancel AND Dispose the
|
|
// replaced stream's CancellationTokenSource. That CTS belongs to the replaced
|
|
// handler's own `using var streamCts`, which is still running and still has to
|
|
// read `streamCts.Token` — so the Dispose raced that read and escaped the RPC as
|
|
// an unhandled ObjectDisposedException. It surfaced only under full-suite load
|
|
// (GrpcStreamIntegrationTests.Pipeline_DuplicateCorrelationId_ReplacesStream) and
|
|
// predates R2: the same Dispose and the same first-token-read relationship existed
|
|
// when the handler still used `ReadAllAsync(streamCts.Token)`.
|
|
//
|
|
// The race is made DETERMINISTIC here by gating the first stream inside its setup
|
|
// window (its _activeStreams entry is registered before Subscribe is called), so
|
|
// the replacement always lands before the first stream reads its token.
|
|
using var gate = new ManualResetEventSlim(false);
|
|
var calls = 0;
|
|
var subscriber = Substitute.For<ISiteStreamSubscriber>();
|
|
subscriber.Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>())
|
|
.Returns(ci =>
|
|
{
|
|
var n = Interlocked.Increment(ref calls);
|
|
if (n == 1)
|
|
gate.Wait(TimeSpan.FromSeconds(15));
|
|
return $"sub-dup-race-{n}";
|
|
});
|
|
|
|
var server = new SiteStreamGrpcServer(subscriber, _logger);
|
|
server.SetReady(Sys);
|
|
|
|
using var cts1 = new CancellationTokenSource();
|
|
var stream1 = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-dup-race"),
|
|
Substitute.For<IServerStreamWriter<SiteStreamEvent>>(),
|
|
CreateMockContext(cts1.Token)));
|
|
|
|
await WaitForConditionAsync(() => server.ActiveStreamCount == 1);
|
|
await WaitForConditionAsync(() => Volatile.Read(ref calls) == 1);
|
|
|
|
using var cts2 = new CancellationTokenSource();
|
|
var stream2 = Task.Run(() => server.SubscribeInstance(
|
|
MakeRequest("corr-dup-race"),
|
|
Substitute.For<IServerStreamWriter<SiteStreamEvent>>(),
|
|
CreateMockContext(cts2.Token)));
|
|
|
|
// The replacement has taken the slot (and cancelled stream 1's CTS) by the time
|
|
// its own Subscribe has been called.
|
|
await WaitForConditionAsync(() => Volatile.Read(ref calls) == 2);
|
|
|
|
gate.Set();
|
|
|
|
// Pre-fix this threw ObjectDisposedException out of the RPC. Post-fix the replaced
|
|
// stream observes a plain cancellation and unwinds through its normal finally.
|
|
await stream1;
|
|
|
|
cts2.Cancel();
|
|
await stream2;
|
|
|
|
Assert.Equal(0, server.ActiveStreamCount);
|
|
}
|
|
|
|
// ── R2: gRPC event batching, and its negotiation ────────────────────────────
|
|
|
|
[Fact]
|
|
public void BatchOptions_AreBoundFromOptions_AndClampDegenerateValues()
|
|
{
|
|
var options = Microsoft.Extensions.Options.Options.Create(new CommunicationOptions());
|
|
var server = new SiteStreamGrpcServer(_subscriber, _logger, options);
|
|
|
|
Assert.Equal(100, server.StreamBatchMaxEvents);
|
|
Assert.Equal(TimeSpan.FromMilliseconds(25), server.StreamBatchWindow);
|
|
|
|
// CommunicationOptionsValidator fails the boot on these, but a host composed
|
|
// without validation must not blow up deep inside a live RPC.
|
|
var degenerate = new SiteStreamGrpcServer(_subscriber, _logger,
|
|
Microsoft.Extensions.Options.Options.Create(new CommunicationOptions
|
|
{
|
|
GrpcStreamBatchMaxEvents = 0,
|
|
GrpcStreamBatchWindow = TimeSpan.FromMilliseconds(-5),
|
|
}));
|
|
Assert.Equal(1, degenerate.StreamBatchMaxEvents);
|
|
Assert.Equal(TimeSpan.Zero, degenerate.StreamBatchWindow);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task UnnegotiatedSubscription_NeverEmitsABatchFrame()
|
|
{
|
|
// OLD-CENTRAL ↔ NEW-SITE skew. proto3 defaults batching_supported to false, which
|
|
// is exactly what a central built before R2 sends. The site must then keep to one
|
|
// event per frame — a Batch frame would arrive at that central as
|
|
// EventOneofCase.None and be silently dropped by its ConvertToDomainEvent.
|
|
var (server, capture, cts, streamTask, relay) =
|
|
await StartCapturingStreamAsync(batchingSupported: false);
|
|
|
|
for (var i = 0; i < 50; i++)
|
|
{
|
|
relay.Tell(new Commons.Messages.Streaming.AttributeValueChanged(
|
|
"Site1.Pump01", "Path", "Attr", i, "Good", DateTimeOffset.UtcNow));
|
|
}
|
|
|
|
await WaitForConditionAsync(() => CountEvents(capture) >= 50, 10_000);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
lock (capture)
|
|
{
|
|
Assert.All(capture, f => Assert.NotEqual(SiteStreamEvent.EventOneofCase.Batch, f.EventCase));
|
|
Assert.Equal(50, capture.Count);
|
|
}
|
|
|
|
GC.KeepAlive(server);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task NegotiatedSubscription_CoalescesABacklogIntoFewerFramesThanEvents()
|
|
{
|
|
// NEW-CENTRAL ↔ NEW-SITE. A burst pushed at the relay faster than the pump drains
|
|
// it must come out in strictly fewer frames than events, with every event
|
|
// preserved in order.
|
|
const int burst = 400;
|
|
var (server, capture, cts, streamTask, relay) =
|
|
await StartCapturingStreamAsync(batchingSupported: true);
|
|
|
|
for (var i = 0; i < burst; i++)
|
|
{
|
|
relay.Tell(new Commons.Messages.Streaming.AttributeValueChanged(
|
|
"Site1.Pump01", "Path", "Attr", i, "Good", DateTimeOffset.UtcNow));
|
|
}
|
|
|
|
await WaitForConditionAsync(() => CountEvents(capture) >= burst, 15_000);
|
|
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
List<SiteStreamEvent> frames;
|
|
lock (capture) { frames = [.. capture]; }
|
|
|
|
Assert.Equal(burst, frames.Sum(CountFrameEvents));
|
|
Assert.True(frames.Count < burst,
|
|
$"batching produced {frames.Count} frames for {burst} events — no coalescing happened");
|
|
Assert.Contains(frames, f => f.EventCase == SiteStreamEvent.EventOneofCase.Batch);
|
|
|
|
// Order is preserved end to end: the values arrive 0..burst-1 exactly once each.
|
|
var values = frames.SelectMany(FlattenAttributeValues).ToArray();
|
|
Assert.Equal(Enumerable.Range(0, burst).Select(i => i.ToString()), values);
|
|
|
|
GC.KeepAlive(server);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task BatchSizeHistogram_IsRecordedOnlyForNegotiatedStreams()
|
|
{
|
|
// scadabridge.site.stream.batch_size rides ScadaBridgeTelemetry.MeterName, which is
|
|
// already in SiteServiceRegistration.ObservedMeters — an unlisted meter exports
|
|
// nothing, silently. Assert the instrument actually fires, and that it does NOT
|
|
// fire on an un-negotiated stream (where it would degenerate into a per-event
|
|
// instrument on the hottest path in the product).
|
|
var measurements = new List<int>();
|
|
using var listener = new MeterListener();
|
|
listener.InstrumentPublished = (instrument, l) =>
|
|
{
|
|
if (instrument.Meter.Name == ScadaBridgeTelemetry.MeterName &&
|
|
instrument.Name == "scadabridge.site.stream.batch_size")
|
|
{
|
|
l.EnableMeasurementEvents(instrument);
|
|
}
|
|
};
|
|
listener.SetMeasurementEventCallback<int>((_, m, _, _) =>
|
|
{
|
|
lock (measurements) { measurements.Add(m); }
|
|
});
|
|
listener.Start();
|
|
|
|
// Un-negotiated: no measurements at all.
|
|
var (_, plainCapture, plainCts, plainTask, plainRelay) =
|
|
await StartCapturingStreamAsync(batchingSupported: false, correlationId: "corr-hist-off");
|
|
plainRelay.Tell(new Commons.Messages.Streaming.AttributeValueChanged(
|
|
"Site1.Pump01", "Path", "Attr", 1, "Good", DateTimeOffset.UtcNow));
|
|
await WaitForConditionAsync(() => CountEvents(plainCapture) >= 1);
|
|
plainCts.Cancel();
|
|
await plainTask;
|
|
|
|
lock (measurements) { Assert.Empty(measurements); }
|
|
|
|
// Negotiated: one measurement per emitted frame, each within the size cap.
|
|
var (_, capture, cts, streamTask, relay) =
|
|
await StartCapturingStreamAsync(batchingSupported: true, correlationId: "corr-hist-on");
|
|
for (var i = 0; i < 20; i++)
|
|
{
|
|
relay.Tell(new Commons.Messages.Streaming.AttributeValueChanged(
|
|
"Site1.Pump01", "Path", "Attr", i, "Good", DateTimeOffset.UtcNow));
|
|
}
|
|
await WaitForConditionAsync(() => CountEvents(capture) >= 20, 10_000);
|
|
cts.Cancel();
|
|
await streamTask;
|
|
|
|
lock (measurements)
|
|
{
|
|
Assert.NotEmpty(measurements);
|
|
Assert.Equal(20, measurements.Sum());
|
|
Assert.All(measurements, m => Assert.InRange(m, 1, SiteStreamGrpcServer.DefaultStreamBatchMaxEvents));
|
|
}
|
|
}
|
|
|
|
/// <summary>Total events carried across all captured frames (unpacking batch frames).</summary>
|
|
private static int CountEvents(List<SiteStreamEvent> capture)
|
|
{
|
|
lock (capture) { return capture.Sum(CountFrameEvents); }
|
|
}
|
|
|
|
private static int CountFrameEvents(SiteStreamEvent frame) =>
|
|
frame.EventCase == SiteStreamEvent.EventOneofCase.Batch ? frame.Batch.Events.Count : 1;
|
|
|
|
private static IEnumerable<string> FlattenAttributeValues(SiteStreamEvent frame)
|
|
{
|
|
if (frame.EventCase == SiteStreamEvent.EventOneofCase.Batch)
|
|
{
|
|
foreach (var inner in frame.Batch.Events)
|
|
yield return inner.AttributeChanged.Value;
|
|
yield break;
|
|
}
|
|
|
|
yield return frame.AttributeChanged.Value;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Starts a SubscribeInstance stream with the given batch negotiation, capturing every
|
|
/// written frame and handing back the relay actor so the test can drive domain events.
|
|
/// </summary>
|
|
private async Task<(SiteStreamGrpcServer Server, List<SiteStreamEvent> Capture,
|
|
CancellationTokenSource Cts, Task StreamTask, IActorRef Relay)>
|
|
StartCapturingStreamAsync(bool batchingSupported, string correlationId = "corr-batch")
|
|
{
|
|
IActorRef? capturedActor = null;
|
|
var subscriber = Substitute.For<ISiteStreamSubscriber>();
|
|
subscriber.Subscribe(Arg.Any<string>(), Arg.Any<IActorRef>())
|
|
.Returns(ci =>
|
|
{
|
|
capturedActor = ci.Arg<IActorRef>();
|
|
return "sub-batch";
|
|
});
|
|
|
|
var server = new SiteStreamGrpcServer(subscriber, _logger,
|
|
Microsoft.Extensions.Options.Options.Create(new CommunicationOptions()));
|
|
server.SetReady(Sys);
|
|
|
|
var capture = new List<SiteStreamEvent>();
|
|
var writer = Substitute.For<IServerStreamWriter<SiteStreamEvent>>();
|
|
writer.WriteAsync(Arg.Any<SiteStreamEvent>(), Arg.Any<CancellationToken>())
|
|
.Returns(Task.CompletedTask)
|
|
.AndDoes(ci =>
|
|
{
|
|
var frame = ci.Arg<SiteStreamEvent>();
|
|
lock (capture) { capture.Add(frame); }
|
|
});
|
|
|
|
var cts = new CancellationTokenSource();
|
|
var context = CreateMockContext(cts.Token);
|
|
|
|
var request = new InstanceStreamRequest
|
|
{
|
|
CorrelationId = correlationId,
|
|
InstanceUniqueName = "Site1.Pump01",
|
|
BatchingSupported = batchingSupported
|
|
};
|
|
|
|
var streamTask = Task.Run(() => server.SubscribeInstance(request, writer, context));
|
|
await WaitForConditionAsync(() => capturedActor != null);
|
|
|
|
return (server, capture, cts, streamTask, capturedActor!);
|
|
}
|
|
}
|