274 lines
11 KiB
C#
274 lines
11 KiB
C#
using Akka.TestKit.Xunit2;
|
|
using Google.Protobuf.WellKnownTypes;
|
|
using Grpc.Core;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
using NSubstitute;
|
|
using NSubstitute.ExceptionExtensions;
|
|
using ZB.MOM.WW.Audit;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Services;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Types.Audit;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
|
|
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
|
|
|
|
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
|
|
|
|
/// <summary>
|
|
/// Tests for <see cref="SiteStreamGrpcServer.PullAuditEvents"/>: the request →
|
|
/// <c>ISiteAuditQueue.ReadPendingSinceAsync</c> → response round-trip, plus the WP2.3
|
|
/// at-least-once contract — rows are retired by the NEXT pull's cursor
|
|
/// (<c>MarkReconciledUpToAsync</c>), never by the act of serving them. The queue is an
|
|
/// NSubstitute stub so the tests never touch SQLite.
|
|
/// </summary>
|
|
public class SiteStreamPullAuditEventsTests : TestKit
|
|
{
|
|
private readonly ISiteStreamSubscriber _subscriber = Substitute.For<ISiteStreamSubscriber>();
|
|
|
|
private SiteStreamGrpcServer CreateServer() =>
|
|
new(_subscriber, NullLogger<SiteStreamGrpcServer>.Instance);
|
|
|
|
private static ServerCallContext NewContext(CancellationToken ct = default)
|
|
{
|
|
var context = Substitute.For<ServerCallContext>();
|
|
context.CancellationToken.Returns(ct);
|
|
return context;
|
|
}
|
|
|
|
// C3 (Task 2.5): canonical ZB.MOM.WW.Audit.AuditEvent via the shared factory.
|
|
// ForwardState is no longer a record field — it is a site-storage-only concern.
|
|
private static AuditEvent NewEvent(DateTime? occurredAt = null) =>
|
|
ScadaBridgeAuditEventFactory.Create(
|
|
channel: AuditChannel.ApiOutbound,
|
|
kind: AuditKind.ApiCall,
|
|
status: AuditStatus.Delivered,
|
|
occurredAtUtc: occurredAt
|
|
?? DateTime.SpecifyKind(new DateTime(2026, 5, 20, 10, 0, 0), DateTimeKind.Utc),
|
|
sourceSiteId: "site-1");
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_NoQueueWired_ReturnsEmptyResponse()
|
|
{
|
|
var server = CreateServer();
|
|
// Intentionally do NOT call SetSiteAuditQueue — simulates a central-only
|
|
// host or a wiring-incomplete startup window.
|
|
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(DateTime.UtcNow.AddMinutes(-5)),
|
|
BatchSize = 100,
|
|
};
|
|
|
|
var response = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Empty(response.Events);
|
|
Assert.False(response.MoreAvailable);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_With5PendingRows_ReturnsAllFiveDtos_AndDoesNotFlipThem()
|
|
{
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
var events = Enumerable.Range(0, 5).Select(_ => NewEvent()).ToList();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)events);
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(DateTime.UtcNow.AddHours(-1)),
|
|
BatchSize = 100, // larger than returned count so MoreAvailable should be false
|
|
};
|
|
|
|
var response = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Equal(5, response.Events.Count);
|
|
Assert.False(response.MoreAvailable); // 5 < 100
|
|
var expectedIds = events.Select(e => e.EventId.ToString()).ToHashSet();
|
|
Assert.True(expectedIds.SetEquals(response.Events.Select(d => d.EventId).ToHashSet()));
|
|
|
|
// AT-LEAST-ONCE: serving rows is NOT proof of receipt. The per-id flip is gone
|
|
// entirely; only a later cursor retires rows.
|
|
await queue.DidNotReceive().MarkReconciledAsync(
|
|
Arg.Any<IReadOnlyList<Guid>>(), Arg.Any<CancellationToken>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_FaultBetweenResponseAndNextPull_ReservesTheSameRows()
|
|
{
|
|
// The failure this closes: central receives the batch, then dies before committing
|
|
// it, so its cursor never advances. Pre-fix the site had already flipped the rows to
|
|
// Reconciled while serving them, and ReadPendingSinceAsync would never return them
|
|
// again — the rows were silently lost. Now the unchanged cursor means no flip, and
|
|
// the identical batch is served again.
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
var since = DateTime.SpecifyKind(new DateTime(2026, 5, 20, 9, 30, 0), DateTimeKind.Utc);
|
|
var events = Enumerable.Range(0, 3).Select(_ => NewEvent()).ToList();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)events);
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(since),
|
|
BatchSize = 100,
|
|
};
|
|
|
|
var first = await server.PullAuditEvents(request, NewContext());
|
|
// …central faults here; it never commits, so it re-pulls with the SAME cursor.
|
|
var second = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Equal(3, first.Events.Count);
|
|
Assert.Equal(
|
|
first.Events.Select(e => e.EventId).ToHashSet(),
|
|
second.Events.Select(e => e.EventId).ToHashSet());
|
|
|
|
// Neither pull retired anything past the (unchanged) cursor: the flip is bounded by
|
|
// the cursor value, so replaying the same cursor can never retire the served rows.
|
|
await queue.Received(2).MarkReconciledUpToAsync(
|
|
since, null, Arg.Any<CancellationToken>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_AdvancedCursor_RetiresEverythingUpToIt_BeforeReading()
|
|
{
|
|
// The cursor central sends back IS the receipt: everything at or before it has been
|
|
// ingested, so those rows are flipped — and flipped BEFORE the read, so they do not
|
|
// consume this batch's budget.
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)Array.Empty<AuditEvent>());
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var cursorTime = DateTime.SpecifyKind(new DateTime(2026, 5, 20, 10, 0, 0), DateTimeKind.Utc);
|
|
var cursorId = Guid.NewGuid().ToString();
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(cursorTime),
|
|
BatchSize = 100,
|
|
AfterId = cursorId,
|
|
};
|
|
|
|
await server.PullAuditEvents(request, NewContext());
|
|
|
|
await queue.Received(1).MarkReconciledUpToAsync(
|
|
cursorTime, cursorId, Arg.Any<CancellationToken>());
|
|
// The keyset cursor is passed straight through to the read as well.
|
|
await queue.Received(1).ReadPendingSinceAsync(
|
|
cursorTime, 100, cursorId, Arg.Any<CancellationToken>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_FirstEverPull_DoesNotFlipAnything()
|
|
{
|
|
// since == MinValue means "from the beginning of recorded history" — central has
|
|
// consumed nothing yet, so there is nothing to retire.
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)Array.Empty<AuditEvent>());
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
await server.PullAuditEvents(new PullAuditEventsRequest { BatchSize = 10 }, NewContext());
|
|
|
|
await queue.DidNotReceive().MarkReconciledUpToAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<string?>(), Arg.Any<CancellationToken>());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_RowsOlderThanSinceUtc_Excluded()
|
|
{
|
|
// The handler delegates the since-utc filter to ReadPendingSinceAsync;
|
|
// this test verifies it passes the request value through verbatim
|
|
// (no clock skew, no off-by-one) and that an empty queue response
|
|
// yields an empty gRPC response.
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
var capturedSince = DateTime.MinValue;
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns(call =>
|
|
{
|
|
capturedSince = call.ArgAt<DateTime>(0);
|
|
return (IReadOnlyList<AuditEvent>)Array.Empty<AuditEvent>();
|
|
});
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var since = DateTime.SpecifyKind(new DateTime(2026, 5, 20, 9, 30, 0), DateTimeKind.Utc);
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(since),
|
|
BatchSize = 50,
|
|
};
|
|
|
|
var response = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Empty(response.Events);
|
|
Assert.False(response.MoreAvailable);
|
|
Assert.Equal(since, capturedSince);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_BatchSize3_Returns3Rows_MoreAvailableTrue()
|
|
{
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
var events = Enumerable.Range(0, 3).Select(_ => NewEvent()).ToList();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)events);
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(DateTime.UtcNow.AddHours(-1)),
|
|
BatchSize = 3,
|
|
};
|
|
|
|
var response = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Equal(3, response.Events.Count);
|
|
// saturated batch → central needs to know to issue a follow-up pull
|
|
Assert.True(response.MoreAvailable);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task PullAuditEvents_MarkReconciledUpToThrows_ResponseStillReturned()
|
|
{
|
|
// The retire step is best-effort — if it fails, the pull must still serve rows.
|
|
// Worst case the same rows are shipped again and central dedups on EventId.
|
|
var queue = Substitute.For<ISiteAuditQueue>();
|
|
var events = Enumerable.Range(0, 2).Select(_ => NewEvent()).ToList();
|
|
queue.ReadPendingSinceAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<int>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.Returns((IReadOnlyList<AuditEvent>)events);
|
|
queue.MarkReconciledUpToAsync(
|
|
Arg.Any<DateTime>(), Arg.Any<string?>(), Arg.Any<CancellationToken>())
|
|
.ThrowsAsync(new InvalidOperationException("SQLite disposed mid-call"));
|
|
|
|
var server = CreateServer();
|
|
server.SetSiteAuditQueue(queue);
|
|
|
|
var request = new PullAuditEventsRequest
|
|
{
|
|
SinceUtc = Timestamp.FromDateTime(DateTime.UtcNow.AddHours(-1)),
|
|
BatchSize = 100,
|
|
};
|
|
|
|
var response = await server.PullAuditEvents(request, NewContext());
|
|
|
|
Assert.Equal(2, response.Events.Count);
|
|
}
|
|
}
|