158 lines
6.1 KiB
C#
158 lines
6.1 KiB
C#
using System.Text.Json;
|
|
using Akka.Actor;
|
|
using Akka.TestKit.Xunit2;
|
|
using ScadaLink.Commons.Messages.Notification;
|
|
using ScadaLink.Commons.Types.Enums;
|
|
|
|
namespace ScadaLink.StoreAndForward.Tests;
|
|
|
|
/// <summary>
|
|
/// Notification Outbox: tests for the site Store-and-Forward notification delivery
|
|
/// handler. "Delivering" a buffered notification means forwarding it to the central
|
|
/// cluster (via the site communication actor) and treating central's
|
|
/// <see cref="NotificationSubmitAck"/> as the delivery outcome.
|
|
/// </summary>
|
|
public class NotificationForwarderTests : TestKit
|
|
{
|
|
private static readonly TimeSpan ForwardTimeout = TimeSpan.FromSeconds(2);
|
|
|
|
/// <summary>
|
|
/// Builds a buffered notification S&F message whose payload matches the shape
|
|
/// produced by the site NotificationDeliveryService enqueue path.
|
|
/// </summary>
|
|
private static StoreAndForwardMessage BufferedNotification(
|
|
string id = "msg-1", string listName = "Operators",
|
|
string subject = "Pump alarm", string message = "Pump 3 tripped",
|
|
string? originInstance = "Plant.Pump3")
|
|
{
|
|
var payload = JsonSerializer.Serialize(new
|
|
{
|
|
ListName = listName,
|
|
Subject = subject,
|
|
Message = message
|
|
});
|
|
return new StoreAndForwardMessage
|
|
{
|
|
Id = id,
|
|
Category = StoreAndForwardCategory.Notification,
|
|
Target = listName,
|
|
PayloadJson = payload,
|
|
OriginInstanceName = originInstance,
|
|
};
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Deliver_ForwardsNotificationSubmitToCentralTarget_AndReturnsTrueOnAccept()
|
|
{
|
|
var centralProbe = CreateTestProbe();
|
|
var forwarder = new NotificationForwarder(
|
|
centralProbe.Ref, "site-7", ForwardTimeout);
|
|
|
|
var msg = BufferedNotification(
|
|
id: "msg-1", listName: "Operators", subject: "Pump alarm",
|
|
message: "Pump 3 tripped", originInstance: "Plant.Pump3");
|
|
|
|
var deliverTask = forwarder.DeliverAsync(msg);
|
|
|
|
// The central target receives a NotificationSubmit whose fields map from the
|
|
// buffered payload; reply Accepted so the handler completes as delivered.
|
|
var submit = centralProbe.ExpectMsg<NotificationSubmit>();
|
|
Assert.Equal("Operators", submit.ListName);
|
|
Assert.Equal("Pump alarm", submit.Subject);
|
|
Assert.Equal("Pump 3 tripped", submit.Body);
|
|
Assert.Equal("site-7", submit.SourceSiteId);
|
|
Assert.Equal("Plant.Pump3", submit.SourceInstanceId);
|
|
centralProbe.Reply(new NotificationSubmitAck(submit.NotificationId, Accepted: true, Error: null));
|
|
|
|
Assert.True(await deliverTask);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Deliver_FallsBackToTarget_WhenPayloadListNameIsEmpty()
|
|
{
|
|
var centralProbe = CreateTestProbe();
|
|
var forwarder = new NotificationForwarder(
|
|
centralProbe.Ref, "site-7", ForwardTimeout);
|
|
|
|
// A buffered payload carrying an empty-string ListName: the empty value must not
|
|
// be forwarded — the forwarder falls back to the S&F message Target instead.
|
|
var payload = JsonSerializer.Serialize(new
|
|
{
|
|
ListName = "",
|
|
Subject = "Pump alarm",
|
|
Message = "Pump 3 tripped"
|
|
});
|
|
var msg = new StoreAndForwardMessage
|
|
{
|
|
Id = "msg-empty-list",
|
|
Category = StoreAndForwardCategory.Notification,
|
|
Target = "Operators",
|
|
PayloadJson = payload,
|
|
OriginInstanceName = "Plant.Pump3",
|
|
};
|
|
|
|
var deliverTask = forwarder.DeliverAsync(msg);
|
|
|
|
var submit = centralProbe.ExpectMsg<NotificationSubmit>();
|
|
Assert.Equal("Operators", submit.ListName);
|
|
centralProbe.Reply(new NotificationSubmitAck(submit.NotificationId, Accepted: true, Error: null));
|
|
|
|
Assert.True(await deliverTask);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Deliver_ThrowsTransient_WhenAckIsNotAccepted()
|
|
{
|
|
var centralProbe = CreateTestProbe();
|
|
var forwarder = new NotificationForwarder(
|
|
centralProbe.Ref, "site-7", ForwardTimeout);
|
|
|
|
var deliverTask = forwarder.DeliverAsync(BufferedNotification());
|
|
|
|
var submit = centralProbe.ExpectMsg<NotificationSubmit>();
|
|
centralProbe.Reply(new NotificationSubmitAck(
|
|
submit.NotificationId, Accepted: false, Error: "central rejected"));
|
|
|
|
// A non-accepted ack is a transient failure — the handler throws so the S&F
|
|
// engine keeps the message buffered and retries the forward.
|
|
await Assert.ThrowsAnyAsync<Exception>(() => deliverTask);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Deliver_ThrowsTransient_WhenNoReplyWithinTimeout()
|
|
{
|
|
// A probe that never replies stands in for central being unreachable.
|
|
var centralProbe = CreateTestProbe();
|
|
var forwarder = new NotificationForwarder(
|
|
centralProbe.Ref, "site-7", TimeSpan.FromMilliseconds(300));
|
|
|
|
// No reply within the timeout → transient failure → throw.
|
|
await Assert.ThrowsAnyAsync<Exception>(() => forwarder.DeliverAsync(BufferedNotification()));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Deliver_UsesStableNotificationId_AcrossRetriesOfSameMessage()
|
|
{
|
|
var centralProbe = CreateTestProbe();
|
|
var forwarder = new NotificationForwarder(
|
|
centralProbe.Ref, "site-7", ForwardTimeout);
|
|
|
|
var buffered = BufferedNotification(id: "stable-msg-id");
|
|
|
|
var first = forwarder.DeliverAsync(buffered);
|
|
var submit1 = centralProbe.ExpectMsg<NotificationSubmit>();
|
|
centralProbe.Reply(new NotificationSubmitAck(submit1.NotificationId, true, null));
|
|
await first;
|
|
|
|
var second = forwarder.DeliverAsync(buffered);
|
|
var submit2 = centralProbe.ExpectMsg<NotificationSubmit>();
|
|
centralProbe.Reply(new NotificationSubmitAck(submit2.NotificationId, true, null));
|
|
await second;
|
|
|
|
// The NotificationId is the central idempotency key — it must be identical for
|
|
// every forward attempt of the same buffered S&F message.
|
|
Assert.Equal(submit1.NotificationId, submit2.NotificationId);
|
|
Assert.Equal("stable-msg-id", submit1.NotificationId);
|
|
}
|
|
}
|