f2efeb37b7
Tasks 5 and 6 of the Phase 2 plan, committed together because their test
fallout is entangled — several fixtures construct both stores.
StoreAndForwardStorage and SiteStorageService now take ILocalDb. Connections
come from ILocalDb.CreateConnection(), which hands out an already-open,
pragma-configured connection carrying the zb_hlc_next() UDF the capture triggers
call; a raw connection would lack the UDF and every write to a replicated table
would fail closed. Deleted with the connection strings: S&F's
EnsureDatabaseDirectoryExists and its per-open busy_timeout pragma, and the site
service's BusyTimeoutFloorSeconds normalization — LocalDb owns all of it now.
DI: AddSiteRuntime's string overload is gone (nothing left to supply), so the
Host calls the no-arg form. ScadaBridge:Database:SiteDbPath and
StoreAndForwardOptions.SqliteDbPath survive only as the migrator's source
locations in Tasks 8/9.
Two things the plan did not anticipate, both worth reading:
1. FOUND A REAL LATENT DEFECT, from Phase 1, now fixed. The plan assumed
directory creation simply moved to LocalDb along with file ownership. It did
not: the LocalDb library never creates the parent directory, and
SqliteLocalDb opens the file eagerly in its constructor — so a missing
directory is a hard boot failure ("SQLite Error 14: unable to open database
file"), not a degraded start. The default site config points at the RELATIVE
path ./data/site-localdb.db, so any site node without a pre-existing data/
directory fails to boot. The docker rig escapes only because its volume mount
happens to create /app/data — a coincidence that would have hidden this until
a bare-metal or fresh deployment. This has been latent since Phase 1 made
LocalDb:Path required; deleting S&F's EnsureDatabaseDirectoryExists here
would have widened it. Re-established the guarantee at the layer that now
owns the path (SiteLocalDbDirectory.Ensure, called before AddZbLocalDb) and
pinned it with SiteLocalDbDirectoryTests. Non-vacuity is not assumed: two
tests written against the wrong assumption failed with exactly this
SQLite Error 14 before the fix existed.
2. Test fallout was ~7x the plan's estimate. The plan named "fixtures" in one
project; the constructor change actually reaches 40 files across 7 test
projects, and most used Mode=Memory;Cache=Shared — which LocalDb has no
equivalent for, so every one had to move to a real temp file. Rather than
copy the Phase 1 TestLocalDb fixture into 7 projects, added a shared
tests/ZB.MOM.WW.ScadaBridge.TestSupport library (not a test project) so the
WAL-sidecar cleanup and the "real, not stubbed" rationale live in one place.
Retargeted rather than deleted, in both directions: the S&F WAL test now asserts
against the LocalDb-backed store (WAL genuinely is LocalDb's job), while the
directory-creation test moved to Host.Tests (that guarantee is NOT LocalDb's).
SiteStorageServiceTests.Initialize_EnablesWalJournalMode got the same treatment.
DeploymentManagerMediumFindingsTests induced a persistence failure via an
unopenable path, which no longer reaches the assertion since the fixture now
throws first; it induces the same failure shape via an uninitialized store.
Verified: full solution build 0 warnings; SiteRuntime 532, Host 318,
AuditLog 355, ExternalSystemGateway 142, HealthMonitoring 97,
StoreAndForward 153 — 1597 passed, 0 failed.
Claude-Session: https://claude.ai/code/session_01BL2Vu1ESDQ9SCN4gVKkdts
267 lines
10 KiB
C#
267 lines
10 KiB
C#
using Microsoft.Extensions.Logging.Abstractions;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
|
|
using ZB.MOM.WW.ScadaBridge.TestSupport;
|
|
|
|
namespace ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests;
|
|
|
|
/// <summary>
|
|
/// StoreAndForward-001: the active node must forward every buffer operation
|
|
/// (add / remove / park) to the standby via the ReplicationService, so a
|
|
/// failover does not lose the buffer.
|
|
/// </summary>
|
|
public class StoreAndForwardReplicationTests : IAsyncLifetime, IDisposable
|
|
{
|
|
private readonly TestLocalDb _localDb;
|
|
private readonly StoreAndForwardStorage _storage;
|
|
private readonly StoreAndForwardService _service;
|
|
private readonly List<ReplicationOperation> _replicated = new();
|
|
|
|
public StoreAndForwardReplicationTests()
|
|
{
|
|
_localDb = TestLocalDb.CreateTemp("ReplTests");
|
|
|
|
_storage = new StoreAndForwardStorage(_localDb.Db, NullLogger<StoreAndForwardStorage>.Instance);
|
|
|
|
var options = new StoreAndForwardOptions
|
|
{
|
|
DefaultRetryInterval = TimeSpan.Zero,
|
|
DefaultMaxRetries = 1,
|
|
RetryTimerInterval = TimeSpan.FromMinutes(10),
|
|
ReplicationEnabled = true,
|
|
};
|
|
|
|
var replication = new ReplicationService(options, NullLogger<ReplicationService>.Instance);
|
|
replication.SetReplicationHandler(op =>
|
|
{
|
|
lock (_replicated) _replicated.Add(op);
|
|
return Task.CompletedTask;
|
|
});
|
|
|
|
_service = new StoreAndForwardService(
|
|
_storage, options, NullLogger<StoreAndForwardService>.Instance, replication);
|
|
}
|
|
|
|
public async Task InitializeAsync() => await _storage.InitializeAsync();
|
|
public Task DisposeAsync() => Task.CompletedTask;
|
|
public void Dispose()
|
|
{
|
|
var path = _localDb.Path;
|
|
_localDb.Dispose();
|
|
TestLocalDb.DeleteFiles(path);
|
|
}
|
|
|
|
/// <summary>Replication is fire-and-forget (Task.Run); poll until the expected ops arrive.</summary>
|
|
private async Task<List<ReplicationOperation>> WaitForReplicationAsync(int count)
|
|
{
|
|
for (var i = 0; i < 100; i++)
|
|
{
|
|
lock (_replicated)
|
|
if (_replicated.Count >= count) return _replicated.ToList();
|
|
await Task.Delay(20);
|
|
}
|
|
lock (_replicated) return _replicated.ToList();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task BufferingAMessage_ReplicatesAnAddOperation()
|
|
{
|
|
// No handler registered → message is buffered → an Add is replicated.
|
|
var result = await _service.EnqueueAsync(
|
|
StoreAndForwardCategory.ExternalSystem, "api", """{}""");
|
|
Assert.True(result.WasBuffered);
|
|
|
|
var ops = await WaitForReplicationAsync(1);
|
|
Assert.Contains(ops, o =>
|
|
o.OperationType == ReplicationOperationType.Add && o.MessageId == result.MessageId);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SuccessfulRetry_ReplicatesARemoveOperation()
|
|
{
|
|
var calls = 0;
|
|
_service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
|
|
_ => ++calls == 1
|
|
? throw new HttpRequestException("transient")
|
|
: Task.FromResult(true));
|
|
|
|
var result = await _service.EnqueueAsync(
|
|
StoreAndForwardCategory.ExternalSystem, "api", """{}""");
|
|
await _service.RetryPendingMessagesAsync();
|
|
|
|
var ops = await WaitForReplicationAsync(2);
|
|
Assert.Contains(ops, o => o.OperationType == ReplicationOperationType.Add);
|
|
Assert.Contains(ops, o =>
|
|
o.OperationType == ReplicationOperationType.Remove && o.MessageId == result.MessageId);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task ParkedMessage_ReplicatesAParkOperation()
|
|
{
|
|
_service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
|
|
_ => throw new HttpRequestException("always fails"));
|
|
|
|
var result = await _service.EnqueueAsync(
|
|
StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
|
|
await _service.RetryPendingMessagesAsync();
|
|
|
|
var ops = await WaitForReplicationAsync(2);
|
|
Assert.Contains(ops, o =>
|
|
o.OperationType == ReplicationOperationType.Park && o.MessageId == result.MessageId);
|
|
}
|
|
|
|
/// <summary>
|
|
/// StoreAndForward-016: an operator discarding a parked message must replicate
|
|
/// a Remove so the standby's copy is also deleted (otherwise the discarded
|
|
/// message reappears in the parked list after a failover).
|
|
/// </summary>
|
|
[Fact]
|
|
public async Task DiscardingAParkedMessage_ReplicatesARemoveOperation()
|
|
{
|
|
_service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
|
|
_ => throw new HttpRequestException("always fails"));
|
|
|
|
var result = await _service.EnqueueAsync(
|
|
StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
|
|
await _service.RetryPendingMessagesAsync(); // -> parked
|
|
await WaitForReplicationAsync(2);
|
|
|
|
var discarded = await _service.DiscardParkedMessageAsync(result.MessageId);
|
|
Assert.True(discarded);
|
|
|
|
var ops = await WaitForReplicationAsync(3);
|
|
Assert.Contains(ops, o =>
|
|
o.OperationType == ReplicationOperationType.Remove && o.MessageId == result.MessageId);
|
|
}
|
|
|
|
/// <summary>
|
|
/// StoreAndForward-016: an operator retrying a parked message must replicate a
|
|
/// Requeue so the standby's copy moves back to Pending (otherwise it stays
|
|
/// Parked on the standby and the operator's retry is lost across a failover).
|
|
/// </summary>
|
|
[Fact]
|
|
public async Task RetryingAParkedMessage_ReplicatesARequeueOperation()
|
|
{
|
|
_service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
|
|
_ => throw new HttpRequestException("always fails"));
|
|
|
|
var result = await _service.EnqueueAsync(
|
|
StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
|
|
await _service.RetryPendingMessagesAsync(); // -> parked
|
|
await WaitForReplicationAsync(2);
|
|
|
|
var retried = await _service.RetryParkedMessageAsync(result.MessageId);
|
|
Assert.True(retried);
|
|
|
|
var ops = await WaitForReplicationAsync(3);
|
|
var requeue = ops.SingleOrDefault(o =>
|
|
o.OperationType == ReplicationOperationType.Requeue && o.MessageId == result.MessageId);
|
|
Assert.NotNull(requeue);
|
|
Assert.NotNull(requeue!.Message);
|
|
Assert.Equal(StoreAndForwardMessageStatus.Pending, requeue.Message!.Status);
|
|
}
|
|
|
|
/// <summary>
|
|
/// StoreAndForward-016: the standby applies a Requeue by moving its row back to
|
|
/// Pending with retry_count = 0, mirroring the active node's local state.
|
|
/// </summary>
|
|
[Fact]
|
|
public async Task ApplyReplicatedOperation_Requeue_MovesStandbyRowBackToPending()
|
|
{
|
|
var replication = new ReplicationService(
|
|
new StoreAndForwardOptions { ReplicationEnabled = true },
|
|
NullLogger<ReplicationService>.Instance);
|
|
|
|
var parked = new StoreAndForwardMessage
|
|
{
|
|
Id = "requeue1",
|
|
Category = StoreAndForwardCategory.ExternalSystem,
|
|
Target = "api",
|
|
PayloadJson = "{}",
|
|
RetryCount = 5,
|
|
MaxRetries = 1,
|
|
RetryIntervalMs = 0,
|
|
CreatedAt = DateTimeOffset.UtcNow,
|
|
Status = StoreAndForwardMessageStatus.Parked,
|
|
};
|
|
await _storage.EnqueueAsync(parked);
|
|
|
|
var requeued = new StoreAndForwardMessage
|
|
{
|
|
Id = parked.Id,
|
|
Category = parked.Category,
|
|
Target = parked.Target,
|
|
PayloadJson = parked.PayloadJson,
|
|
RetryCount = 0,
|
|
MaxRetries = parked.MaxRetries,
|
|
RetryIntervalMs = parked.RetryIntervalMs,
|
|
CreatedAt = parked.CreatedAt,
|
|
Status = StoreAndForwardMessageStatus.Pending,
|
|
};
|
|
await replication.ApplyReplicatedOperationAsync(
|
|
new ReplicationOperation(ReplicationOperationType.Requeue, parked.Id, requeued),
|
|
_storage);
|
|
|
|
var row = await _storage.GetMessageByIdAsync(parked.Id);
|
|
Assert.NotNull(row);
|
|
Assert.Equal(StoreAndForwardMessageStatus.Pending, row!.Status);
|
|
Assert.Equal(0, row.RetryCount);
|
|
}
|
|
|
|
private static ReplicationService NewReplicationService() =>
|
|
new(new StoreAndForwardOptions { ReplicationEnabled = true },
|
|
NullLogger<ReplicationService>.Instance);
|
|
|
|
private static StoreAndForwardMessage NewMessage(string id) => new()
|
|
{
|
|
Id = id,
|
|
Category = StoreAndForwardCategory.ExternalSystem,
|
|
Target = "api",
|
|
PayloadJson = "{}",
|
|
RetryCount = 0,
|
|
MaxRetries = 1,
|
|
RetryIntervalMs = 0,
|
|
CreatedAt = DateTimeOffset.UtcNow,
|
|
Status = StoreAndForwardMessageStatus.Pending,
|
|
};
|
|
|
|
/// <summary>
|
|
/// A Park (or Requeue) whose original Add was lost — e.g. the Add's
|
|
/// fire-and-forget replication dropped — must still materialise the row on the
|
|
/// standby. The full message rides in the operation, so an upsert self-heals;
|
|
/// a blind UPDATE would affect 0 rows and the row would be gone forever.
|
|
/// </summary>
|
|
[Fact]
|
|
public async Task ApplyReplicatedPark_WhenAddWasLost_MaterializesTheParkedRow()
|
|
{
|
|
var service = NewReplicationService();
|
|
var msg = NewMessage("lost-add-1"); // never Added on this (standby) storage
|
|
msg.Status = StoreAndForwardMessageStatus.Parked;
|
|
|
|
await service.ApplyReplicatedOperationAsync(
|
|
new ReplicationOperation(ReplicationOperationType.Park, msg.Id, msg), _storage);
|
|
|
|
var row = await _storage.GetMessageByIdAsync("lost-add-1");
|
|
Assert.NotNull(row); // pre-fix: blind UPDATE affected 0 rows, row is gone forever
|
|
Assert.Equal(StoreAndForwardMessageStatus.Parked, row!.Status);
|
|
}
|
|
|
|
/// <summary>
|
|
/// A duplicate Add (e.g. after Task 21's anti-entropy resync re-issues an Add
|
|
/// the standby already holds) must not violate the PK — the upsert applies
|
|
/// newest-wins instead of throwing.
|
|
/// </summary>
|
|
[Fact]
|
|
public async Task ApplyReplicatedAdd_Twice_IsIdempotent_NewestWins()
|
|
{
|
|
var service = NewReplicationService();
|
|
var msg = NewMessage("dup-add");
|
|
await service.ApplyReplicatedOperationAsync(
|
|
new ReplicationOperation(ReplicationOperationType.Add, msg.Id, msg), _storage);
|
|
msg.RetryCount = 3;
|
|
await service.ApplyReplicatedOperationAsync( // pre-fix: SqliteException PK violation
|
|
new ReplicationOperation(ReplicationOperationType.Add, msg.Id, msg), _storage);
|
|
|
|
Assert.Equal(3, (await _storage.GetMessageByIdAsync("dup-add"))!.RetryCount);
|
|
}
|
|
}
|