15ef1f32a8
I dismissed this finding from the code review as a false positive, reasoning that Akka's ActorBase caches Self in a field and that three passing anti-entropy tests traverse the path. Both premises were wrong. Self resolves through Context, which is [ThreadStatic], and throws NotSupportedException once a continuation resumes on a thread-pool thread. The tests passed because a local SQLite store usually completes await SYNCHRONOUSLY, so the continuation stayed on the mailbox thread and the context was still intact. Correctness therefore depended on store latency and thread-pool timing: green here, broken under a slower or contended store, with the only symptom a per-peer warning every announce interval while nodes silently stopped converging. The live-broadcast fast path masked it further — only the anti-entropy repair path was dead. Captures self on the actor thread and passes it in. Adds GenuinelyAsyncSecretStore to force the async path, a regression test that fails on the unfixed code (20s timeout) and passes in 3s after, and ActorContextAfterAwaitTests pinning the underlying Akka behaviour so the wrong assumption cannot be made again. Audited every remaining Self/Sender access in the actor. 170 pass offline / 184 with the live SQL suite / 1 skip / 0 warnings. Claude-Session: https://claude.ai/code/session_01BL2Vu1ESDQ9SCN4gVKkdts
384 lines
16 KiB
C#
384 lines
16 KiB
C#
using Akka.Actor;
|
|
using Akka.Cluster;
|
|
using Akka.Configuration;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
using ZB.MOM.WW.Secrets.Abstractions;
|
|
using ZB.MOM.WW.Secrets.Replication;
|
|
using ZB.MOM.WW.Secrets.Replicator.AkkaDotNet.Tests.Fakes;
|
|
using ZB.MOM.WW.Secrets.Sqlite;
|
|
|
|
namespace ZB.MOM.WW.Secrets.Replicator.AkkaDotNet.Tests;
|
|
|
|
/// <summary>
|
|
/// Two real cluster nodes, two real SQLite stores, real remoting over loopback — the gate the G-7
|
|
/// design called for, run in-process so it belongs to the ordinary offline suite instead of needing
|
|
/// a deployed rig.
|
|
/// </summary>
|
|
/// <remarks>
|
|
/// Nothing here is faked: each node runs its own <see cref="ActorSystem"/> with Akka.Remote on its
|
|
/// own port, they form a cluster, and rows travel through the actual serializer and the actual
|
|
/// distributed pub/sub mediator. A mock-based test of this behaviour would prove almost nothing —
|
|
/// the failure modes worth catching (self-echo, sender loss across the mediator, tombstones that do
|
|
/// not propagate, a partition that never re-converges) only appear when messages really cross nodes.
|
|
/// </remarks>
|
|
public sealed class TwoNodeClusterReplicationTests : IAsyncLifetime
|
|
{
|
|
private readonly List<string> _dbPaths = [];
|
|
private ActorSystem _systemA = null!;
|
|
private ActorSystem _systemB = null!;
|
|
private SqliteSecretStore _storeA = null!;
|
|
private SqliteSecretStore _storeB = null!;
|
|
private ISecretStore _writableA = null!;
|
|
private IActorRef _replicatorB = null!;
|
|
|
|
// Short so anti-entropy assertions do not dominate the suite's runtime.
|
|
private static readonly TimeSpan AnnounceInterval = TimeSpan.FromMilliseconds(300);
|
|
private static readonly TimeSpan Patience = TimeSpan.FromSeconds(20);
|
|
|
|
public async Task InitializeAsync()
|
|
{
|
|
_storeA = CreateStore();
|
|
_storeB = CreateStore();
|
|
|
|
// Port 0 lets the OS assign free ports, so parallel test runs cannot collide.
|
|
_systemA = ActorSystem.Create("zb-secrets-cluster", ClusterConfig(port: 0));
|
|
int portA = ClusterPort(_systemA);
|
|
|
|
// Node A seeds itself; node B joins it.
|
|
_systemA.Dispose();
|
|
_systemA = ActorSystem.Create("zb-secrets-cluster", ClusterConfig(portA, seedPort: portA));
|
|
_systemB = ActorSystem.Create("zb-secrets-cluster", ClusterConfig(port: 0, seedPort: portA));
|
|
|
|
await AwaitClusterUpAsync(_systemA, expectedMembers: 2);
|
|
await AwaitClusterUpAsync(_systemB, expectedMembers: 2);
|
|
|
|
IActorRef replicatorA = _systemA.ActorOf(
|
|
SecretReplicationActor.Props(_storeA, null, AnnounceInterval), "zb-secret-replication");
|
|
_replicatorB = _systemB.ActorOf(
|
|
SecretReplicationActor.Props(_storeB, null, AnnounceInterval), "zb-secret-replication");
|
|
|
|
_writableA = new ReplicatingSecretStore(
|
|
_storeA, new AkkaSecretReplicator(replicatorA), NullLogger<ReplicatingSecretStore>.Instance);
|
|
|
|
// The mediator gossips subscriptions between nodes; publishing before both sides are
|
|
// registered would silently drop the message.
|
|
await AwaitBothSubscribedAsync();
|
|
}
|
|
|
|
private SqliteSecretStore CreateStore()
|
|
{
|
|
string path = Path.Combine(Path.GetTempPath(), $"zb-akka-{Guid.NewGuid():N}.db");
|
|
_dbPaths.Add(path);
|
|
var factory = new SecretsSqliteConnectionFactory(path);
|
|
new SqliteSecretsStoreMigrator(factory).MigrateAsync(CancellationToken.None).GetAwaiter().GetResult();
|
|
return new SqliteSecretStore(factory);
|
|
}
|
|
|
|
private static Config ClusterConfig(int port, int? seedPort = null)
|
|
{
|
|
string seeds = seedPort is null
|
|
? "[]"
|
|
: $"[\"akka.tcp://zb-secrets-cluster@127.0.0.1:{seedPort}\"]";
|
|
|
|
return ConfigurationFactory.ParseString($$"""
|
|
akka {
|
|
loglevel = WARNING
|
|
actor.provider = cluster
|
|
remote.dot-netty.tcp {
|
|
hostname = "127.0.0.1"
|
|
public-hostname = "127.0.0.1"
|
|
port = {{port}}
|
|
}
|
|
cluster {
|
|
seed-nodes = {{seeds}}
|
|
downing-provider-class = "Akka.Cluster.SBR.SplitBrainResolverProvider"
|
|
}
|
|
}
|
|
""").WithFallback(AkkaSecretsReplication.SerializationConfig);
|
|
}
|
|
|
|
private static int ClusterPort(ActorSystem system) =>
|
|
Cluster.Get(system).SelfAddress.Port
|
|
?? throw new InvalidOperationException("Cluster address has no port.");
|
|
|
|
private static async Task AwaitClusterUpAsync(ActorSystem system, int expectedMembers)
|
|
{
|
|
Cluster cluster = Cluster.Get(system);
|
|
DateTime deadline = DateTime.UtcNow + Patience;
|
|
|
|
while (DateTime.UtcNow < deadline)
|
|
{
|
|
if (cluster.State.Members.Count(m => m.Status == MemberStatus.Up) >= expectedMembers)
|
|
{
|
|
return;
|
|
}
|
|
|
|
await Task.Delay(100);
|
|
}
|
|
|
|
throw new TimeoutException(
|
|
$"Cluster did not reach {expectedMembers} Up members within {Patience}.");
|
|
}
|
|
|
|
// Writes a throwaway secret on A and waits for it to land on B. That round trip is the only
|
|
// reliable signal that pub/sub subscriptions have gossiped across both nodes.
|
|
private async Task AwaitBothSubscribedAsync()
|
|
{
|
|
var probe = new SecretName("warmup/probe");
|
|
DateTime deadline = DateTime.UtcNow + Patience;
|
|
|
|
while (DateTime.UtcNow < deadline)
|
|
{
|
|
await _writableA.UpsertAsync(Row("warmup/probe", 0x01), CancellationToken.None);
|
|
|
|
if (await WaitForAsync(_storeB, probe, _ => true, TimeSpan.FromMilliseconds(500)))
|
|
{
|
|
await _storeA.DeleteAsync(probe, "warmup", CancellationToken.None);
|
|
await _storeB.DeleteAsync(probe, "warmup", CancellationToken.None);
|
|
return;
|
|
}
|
|
}
|
|
|
|
throw new TimeoutException("Nodes did not establish pub/sub subscriptions in time.");
|
|
}
|
|
|
|
private static StoredSecret Row(string name, byte marker) => new()
|
|
{
|
|
Name = new SecretName(name),
|
|
ContentType = SecretContentType.Text,
|
|
Ciphertext = [marker],
|
|
Nonce = [1],
|
|
Tag = [2],
|
|
WrappedDek = [3],
|
|
WrapNonce = [4],
|
|
WrapTag = [5],
|
|
KekId = "sha256:shared-kek",
|
|
Revision = 0,
|
|
CreatedUtc = DateTimeOffset.UtcNow,
|
|
UpdatedUtc = DateTimeOffset.UtcNow,
|
|
};
|
|
|
|
private static async Task<bool> WaitForAsync(
|
|
ISecretStore store, SecretName name, Func<StoredSecret, bool> predicate, TimeSpan? timeout = null)
|
|
{
|
|
DateTime deadline = DateTime.UtcNow + (timeout ?? Patience);
|
|
|
|
while (DateTime.UtcNow < deadline)
|
|
{
|
|
StoredSecret? row = await store.GetAsync(name, CancellationToken.None);
|
|
|
|
if (row is not null && predicate(row))
|
|
{
|
|
return true;
|
|
}
|
|
|
|
await Task.Delay(50);
|
|
}
|
|
|
|
return false;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_secret_written_on_node_A_becomes_resolvable_on_node_B()
|
|
{
|
|
var name = new SecretName("db/password");
|
|
await _writableA.UpsertAsync(Row("db/password", 0xAA), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeB, name, r => r.Ciphertext[0] == 0xAA),
|
|
"Node B never received the secret written on node A.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task The_ciphertext_arrives_byte_identical()
|
|
{
|
|
var name = new SecretName("app/exact");
|
|
await _writableA.UpsertAsync(Row("app/exact", 0x7F), CancellationToken.None);
|
|
|
|
Assert.True(await WaitForAsync(_storeB, name, _ => true));
|
|
|
|
StoredSecret onA = (await _storeA.GetAsync(name, CancellationToken.None))!;
|
|
StoredSecret onB = (await _storeB.GetAsync(name, CancellationToken.None))!;
|
|
|
|
Assert.Equal(onA.Ciphertext, onB.Ciphertext);
|
|
Assert.Equal(onA.WrappedDek, onB.WrappedDek);
|
|
Assert.Equal(onA.KekId, onB.KekId);
|
|
// Verbatim: the peer keeps the originating node's revision and timestamp, or the row would
|
|
// look newer on B than on A and bounce back.
|
|
Assert.Equal(onA.Revision, onB.Revision);
|
|
Assert.Equal(onA.UpdatedUtc, onB.UpdatedUtc);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_delete_on_node_A_tombstones_the_row_on_node_B()
|
|
{
|
|
var name = new SecretName("app/doomed");
|
|
await _writableA.UpsertAsync(Row("app/doomed", 0xBB), CancellationToken.None);
|
|
Assert.True(await WaitForAsync(_storeB, name, _ => true));
|
|
|
|
await _writableA.DeleteAsync(name, "operator", CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeB, name, r => r.IsDeleted),
|
|
"The delete never propagated to node B.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task An_update_overwrites_the_earlier_value_on_the_peer()
|
|
{
|
|
var name = new SecretName("app/rotating");
|
|
await _writableA.UpsertAsync(Row("app/rotating", 0x01), CancellationToken.None);
|
|
Assert.True(await WaitForAsync(_storeB, name, r => r.Ciphertext[0] == 0x01));
|
|
|
|
await _writableA.UpsertAsync(Row("app/rotating", 0x02), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeB, name, r => r.Ciphertext[0] == 0x02 && r.Revision == 1),
|
|
"Node B did not converge on the updated value.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Anti_entropy_repairs_a_write_that_never_got_broadcast()
|
|
{
|
|
// Simulates a dropped publish — the write goes straight to node A's store, bypassing the
|
|
// replicator entirely, so the ONLY route to node B is the periodic manifest exchange. This is
|
|
// the case that makes best-effort publishing safe.
|
|
var name = new SecretName("app/missed-broadcast");
|
|
await _storeA.UpsertAsync(Row("app/missed-broadcast", 0xCC), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeB, name, r => r.Ciphertext[0] == 0xCC),
|
|
"Anti-entropy failed to repair a write that was never broadcast.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Anti_entropy_pulls_in_the_other_direction_too()
|
|
{
|
|
// The mirror case: node B holds a row node A has never heard of. A cannot request a name it
|
|
// does not know exists, so this only converges if the exchange pushes as well as pulls.
|
|
var name = new SecretName("app/born-on-b");
|
|
await _storeB.UpsertAsync(Row("app/born-on-b", 0xDD), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeA, name, r => r.Ciphertext[0] == 0xDD),
|
|
"Node A never learned about a secret that originated on node B.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Anti_entropy_still_works_when_the_store_is_genuinely_asynchronous()
|
|
{
|
|
// Regression guard for a bug the rest of this suite could not see. A local SQLite store
|
|
// usually completes `await` synchronously, so continuations stayed on the actor's mailbox
|
|
// thread and code that touched actor context after an await worked by accident. Wrapping the
|
|
// store so every call really yields makes that illegal access deterministic — this test
|
|
// fails outright if anti-entropy touches actor context across an await.
|
|
SqliteSecretStore backing = CreateStore();
|
|
var slowStore = new GenuinelyAsyncSecretStore(backing);
|
|
|
|
_systemB.ActorOf(
|
|
SecretReplicationActor.Props(slowStore, null, AnnounceInterval), "slow-node");
|
|
|
|
await _writableA.UpsertAsync(Row("slow/inbound", 0x42), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(slowStore, new SecretName("slow/inbound"), r => r.Ciphertext[0] == 0x42),
|
|
"A node whose store is genuinely async never received the secret.");
|
|
|
|
// And the other direction, which can ONLY converge through the anti-entropy exchange.
|
|
await backing.UpsertAsync(Row("slow/outbound", 0x43), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeA, new SecretName("slow/outbound"), r => r.Ciphertext[0] == 0x43),
|
|
"Node A never learned about a secret held by a node whose store is genuinely async.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_node_that_joins_late_catches_up_on_everything()
|
|
{
|
|
// Partition-heal / cold-start: the whole reason to choose this transport over a shared store.
|
|
await _writableA.UpsertAsync(Row("bulk/one", 0x11), CancellationToken.None);
|
|
await _writableA.UpsertAsync(Row("bulk/two", 0x22), CancellationToken.None);
|
|
await _writableA.UpsertAsync(Row("bulk/three", 0x33), CancellationToken.None);
|
|
|
|
SqliteSecretStore lateStore = CreateStore();
|
|
_systemB.ActorOf(
|
|
SecretReplicationActor.Props(lateStore, null, AnnounceInterval), "late-joiner");
|
|
|
|
Assert.True(await WaitForAsync(lateStore, new SecretName("bulk/one"), r => r.Ciphertext[0] == 0x11));
|
|
Assert.True(await WaitForAsync(lateStore, new SecretName("bulk/two"), r => r.Ciphertext[0] == 0x22));
|
|
Assert.True(await WaitForAsync(lateStore, new SecretName("bulk/three"), r => r.Ciphertext[0] == 0x33));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_broadcast_round_trip_does_not_drift_the_originating_nodes_revision()
|
|
{
|
|
// Note on what this does and does NOT prove: the revision would also stay put if the self-echo
|
|
// filter were removed, because re-applying an identical row ties on (updated_utc, revision)
|
|
// and last-writer-wins rejects a tie. So this covers the drift symptom via defence in depth,
|
|
// not the filter itself. The filter is what stops the echo being forwarded onward, which
|
|
// needs three nodes to observe — worth adding if a third node ever joins this rig.
|
|
var name = new SecretName("app/no-echo");
|
|
await _writableA.UpsertAsync(Row("app/no-echo", 0xEE), CancellationToken.None);
|
|
Assert.True(await WaitForAsync(_storeB, name, _ => true));
|
|
|
|
await Task.Delay(AnnounceInterval * 4);
|
|
|
|
StoredSecret onA = (await _storeA.GetAsync(name, CancellationToken.None))!;
|
|
Assert.Equal(0, onA.Revision);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_malformed_row_from_a_peer_is_discarded_without_stalling_replication()
|
|
{
|
|
// Peer input is untrusted: a path-traversing name must be dropped, and — critically — the
|
|
// actor must keep serving everything else afterwards.
|
|
_replicatorB.Tell(new Protocol.SecretRowsMessage(
|
|
[
|
|
new Protocol.SecretRowDto
|
|
{
|
|
Name = "../../etc/passwd",
|
|
ContentType = "Text",
|
|
Ciphertext = [1], Nonce = [1], Tag = [1],
|
|
WrappedDek = [1], WrapNonce = [1], WrapTag = [1],
|
|
KekId = "sha256:x",
|
|
Revision = 0,
|
|
CreatedUtc = DateTimeOffset.UtcNow.ToString("O"),
|
|
UpdatedUtc = DateTimeOffset.UtcNow.ToString("O"),
|
|
},
|
|
]));
|
|
|
|
var name = new SecretName("app/after-malformed");
|
|
await _writableA.UpsertAsync(Row("app/after-malformed", 0x5A), CancellationToken.None);
|
|
|
|
Assert.True(
|
|
await WaitForAsync(_storeB, name, r => r.Ciphertext[0] == 0x5A),
|
|
"Replication stalled after a malformed row was received.");
|
|
}
|
|
|
|
public async Task DisposeAsync()
|
|
{
|
|
await CoordinatedShutdown.Get(_systemA).Run(CoordinatedShutdown.ClrExitReason.Instance);
|
|
await CoordinatedShutdown.Get(_systemB).Run(CoordinatedShutdown.ClrExitReason.Instance);
|
|
|
|
Microsoft.Data.Sqlite.SqliteConnection.ClearAllPools();
|
|
|
|
foreach (string db in _dbPaths)
|
|
{
|
|
foreach (string path in new[] { db, db + "-wal", db + "-shm" })
|
|
{
|
|
try
|
|
{
|
|
if (File.Exists(path))
|
|
{
|
|
File.Delete(path);
|
|
}
|
|
}
|
|
catch (IOException)
|
|
{
|
|
// Best-effort temp cleanup.
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|