feat(cluster): simultaneous-cold-start split-brain guard (#33)

Port OtOpcUa's bootstrap guard (lmxopcua d1dac87f) to ScadaBridge's
BuildHocon bootstrap. Both pair nodes are self-first seeds, so a true
simultaneous cold start (shared site power event) races FirstSeedNodeProcess
on both and forms two 1-node clusters that never merge.

Opt-in dark switch ScadaBridge:Cluster:BootstrapGuard:Enabled (default off,
guard-off behavior byte-identical). When on: BuildHocon emits an empty seed
list so Akka does not auto-join, and ClusterBootstrapCoordinator (IHostedService,
registered in both the Central and Site composition roots) picks the join order
from ClusterBootstrapGuard's pure decision core — the lower canonical host:port
is the founder (self-first, forms immediately); the higher node TCP-probes the
founder up to PartnerProbeSeconds and joins peer-first if reachable, else
self-first (cold-start-alone preserved). Decides before a single JoinSeedNodes,
never re-forms mid-handshake.

Review notes carried over: case-insensitive founder tie-break; fail-fast
validation of probe timings when enabled; higher-node-cold-start-alone covered
by a real-ActorSystem test. 21 unit + 5 real-cluster tests (incl. the headline
both-cold-start-together-form-one-cluster).
This commit is contained in:
Joseph Doherty
2026-07-24 08:55:48 -04:00
parent bb138d5254
commit 248676ed16
10 changed files with 895 additions and 2 deletions
@@ -261,8 +261,16 @@ public class AkkaHostedService : IHostedService
TimeSpan transportHeartbeat,
TimeSpan transportFailure)
{
var seedNodesStr = string.Join(",",
clusterOptions.SeedNodes.Select(QuoteHocon));
// Simultaneous-cold-start split-brain guard (Gitea #33): when enabled the node must NOT
// auto-join from its config seeds — Akka's FirstSeedNodeProcess is exactly what races the
// peer and forms two 1-node clusters. Emit an EMPTY seed list so the ActorSystem comes up
// unjoined, and ClusterBootstrapCoordinator issues the single reachability-gated
// JoinSeedNodes (founder self-first / joiner peer-first). Default OFF, so guard-off configs
// render byte-identical to before.
var effectiveSeeds = clusterOptions.BootstrapGuard.Enabled
? Enumerable.Empty<string>()
: clusterOptions.SeedNodes;
var seedNodesStr = string.Join(",", effectiveSeeds.Select(QuoteHocon));
var rolesStr = string.Join(",", roles.Select(QuoteHocon));
// auto-down (default): AutoDowning provider — the leader among the reachable
@@ -0,0 +1,211 @@
using System.Diagnostics;
using System.Net.Sockets;
using Akka.Actor;
using Akka.Cluster;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.ClusterInfrastructure;
namespace ZB.MOM.WW.ScadaBridge.Host.Actors;
/// <summary>
/// Runs the simultaneous-cold-start split-brain guard (Gitea #33): when
/// <c>ScadaBridge:Cluster:BootstrapGuard:Enabled</c>, the node starts with NO config seed nodes
/// (<c>AkkaHostedService.BuildHocon</c> emits an empty seed list), and this coordinator picks the join
/// order and issues the single <see cref="Akka.Cluster.Cluster.JoinSeedNodes"/> once the ActorSystem is
/// up. See <see cref="ClusterBootstrapGuard"/> for the decision logic and the split-brain rationale.
/// </summary>
/// <remarks>
/// <para>
/// The founder (lower canonical address) joins its self-first order immediately. The higher
/// node probes its partner's Akka endpoint (a plain TCP connect — the partner has bound its
/// port well before it forms) up to <see cref="ClusterBootstrapGuardOptions.PartnerProbeSeconds"/>:
/// reachable ⇒ peer-first (join the founder); unreachable ⇒ self-first (partner is dead, form
/// alone). The join runs on a background task so it never blocks host startup, and it is issued
/// exactly once — the coordinator never re-forms a node mid-handshake, the failure mode of the
/// rejected self-form-watchdog design (see <c>SelfFirstSeedBootstrapTests</c>).
/// </para>
/// <para>
/// Registered unconditionally in both the Central and Site composition roots but a no-op unless
/// the guard is enabled, so guard-off deployments keep Akka's config-driven self-first
/// auto-join byte-identical.
/// </para>
/// </remarks>
public sealed class ClusterBootstrapCoordinator : IHostedService
{
private readonly Func<ActorSystem> _system;
private readonly ClusterOptions _clusterOptions;
private readonly NodeOptions _nodeOptions;
private readonly ILogger<ClusterBootstrapCoordinator> _log;
private readonly CancellationTokenSource _cts = new();
private Task? _joinTask;
/// <summary>Creates the coordinator.</summary>
/// <param name="system">Lazy accessor for the node's ActorSystem (resolved after the Akka hosted service creates it).</param>
/// <param name="clusterOptions">The bound cluster options (seed list + guard settings).</param>
/// <param name="nodeOptions">The bound node identity (advertised hostname + remoting port).</param>
/// <param name="log">Logger for the bootstrap decision.</param>
public ClusterBootstrapCoordinator(
Func<ActorSystem> system,
IOptions<ClusterOptions> clusterOptions,
IOptions<NodeOptions> nodeOptions,
ILogger<ClusterBootstrapCoordinator> log)
{
_system = system ?? throw new ArgumentNullException(nameof(system));
_clusterOptions = clusterOptions?.Value ?? throw new ArgumentNullException(nameof(clusterOptions));
_nodeOptions = nodeOptions?.Value ?? throw new ArgumentNullException(nameof(nodeOptions));
_log = log ?? throw new ArgumentNullException(nameof(log));
}
/// <inheritdoc/>
public Task StartAsync(CancellationToken cancellationToken)
{
if (!_clusterOptions.BootstrapGuard.Enabled)
return Task.CompletedTask; // dark switch off — Akka auto-joined from config seeds already.
// Fire-and-forget so a higher node with a dead partner (which waits the full probe window)
// never blocks host startup. Exceptions are logged; a failed join leaves the node unjoined,
// which is visible and recoverable, never a silent split.
_joinTask = Task.Run(() => RunAsync(_cts.Token), _cts.Token);
return Task.CompletedTask;
}
/// <inheritdoc/>
public async Task StopAsync(CancellationToken cancellationToken)
{
await _cts.CancelAsync().ConfigureAwait(false);
if (_joinTask is not null)
{
try { await _joinTask.ConfigureAwait(false); }
catch (OperationCanceledException) { /* expected on shutdown */ }
}
}
private async Task RunAsync(CancellationToken ct)
{
try
{
var seeds = _clusterOptions.SeedNodes ?? new List<string>();
var role = ClusterBootstrapGuard.Analyze(seeds, _nodeOptions.NodeHostname, _nodeOptions.RemotingPort);
string[] order;
var committedPeerFirst = false;
if (!role.Applies)
{
// Not a pair seed (single-node install, self absent, malformed, or 3+ seeds): join the
// configured seeds unchanged — the guard arbitrates only the 2-node pair race.
order = seeds.ToArray();
_log.LogInformation(
"Bootstrap guard: not a pair seed ({SeedCount} seed(s)); joining configured seeds unchanged.",
seeds.Count);
}
else if (role.IsFounder)
{
order = role.SelfFirstOrder;
_log.LogInformation(
"Bootstrap guard: this node is the preferred founder (lower address); joining self-first, forming immediately if no peer answers.");
}
else
{
var reachable = await ProbePartnerAsync(role.PartnerHost!, role.PartnerPort, ct).ConfigureAwait(false);
order = ClusterBootstrapGuard.HigherNodeOrder(role, reachable);
committedPeerFirst = reachable;
_log.LogInformation(
"Bootstrap guard: this node is the higher address; partner {PartnerHost}:{PartnerPort} was {Reachability} within the probe window; joining {Order}.",
role.PartnerHost, role.PartnerPort, reachable ? "REACHABLE" : "UNREACHABLE",
reachable ? "peer-first (join the founder)" : "self-first (partner down, forming alone)");
}
if (order.Length == 0)
{
_log.LogWarning("Bootstrap guard: no seed nodes to join; node stays unjoined until a seed is configured.");
return;
}
var cluster = Akka.Cluster.Cluster.Get(_system());
var addresses = order.Select(Address.Parse).ToArray();
cluster.JoinSeedNodes(addresses);
// Peer-first is a one-way commitment: JoinSeedNodeProcess never self-forms. If the founder
// died in the probe→join window, this node hangs unjoined forever. We do NOT re-form here
// (the rejected mid-handshake failure mode) — we make the hang operator-visible so a
// restart, which re-runs the guard against the now-dead founder and self-forms, recovers
// it. Nothing is done on the founder / self-first paths, which always self-form on their own.
if (committedPeerFirst)
await WarnIfNotUpAsync(cluster, ct).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
// Host shutting down before the join completed — nothing to do.
}
catch (Exception ex)
{
_log.LogError(ex, "Bootstrap guard: failed to issue JoinSeedNodes; node remains unjoined.");
}
}
/// <summary>
/// After committing peer-first, waits a bounded time for this node to reach <c>Up</c> and logs a
/// clear warning if it does not — the founder must have died in the probe→join window, leaving this
/// node hung in <c>JoinSeedNodeProcess</c>. Read-only: it never re-forms or re-joins (that is the
/// rejected mid-handshake failure mode); it only surfaces the condition so an operator restarts the node.
/// </summary>
private async Task WarnIfNotUpAsync(Akka.Cluster.Cluster cluster, CancellationToken ct)
{
// Give the join a generous grace: the founder's own self-form (seed-node-timeout) plus this
// node's join round-trip. Two probe windows is comfortably beyond both.
var grace = TimeSpan.FromSeconds(Math.Max(10, _clusterOptions.BootstrapGuard.PartnerProbeSeconds * 2));
var sw = Stopwatch.StartNew();
while (!ct.IsCancellationRequested && sw.Elapsed < grace)
{
if (cluster.SelfMember.Status == MemberStatus.Up) return; // joined — all good
try { await Task.Delay(TimeSpan.FromSeconds(1), ct).ConfigureAwait(false); }
catch (OperationCanceledException) { return; }
}
if (!ct.IsCancellationRequested && cluster.SelfMember.Status != MemberStatus.Up)
_log.LogWarning(
"Bootstrap guard: committed peer-first but this node is still {Status} after {Grace}s — its founder likely died in the probe→join window. This node will NOT self-form on its own (peer-first never does). RESTART it to recover: the guard will re-run, find the founder down, and form alone.",
cluster.SelfMember.Status, (int)grace.TotalSeconds);
}
/// <summary>
/// Polls the partner's Akka endpoint with a plain TCP connect until it is reachable or the probe
/// window elapses. Reachable-at-TCP is a sound proxy for "the partner is coming up": a node binds
/// its Akka port near the start of startup, well before it forms or joins a cluster.
/// </summary>
private async Task<bool> ProbePartnerAsync(string host, int port, CancellationToken ct)
{
var window = TimeSpan.FromSeconds(Math.Max(0, _clusterOptions.BootstrapGuard.PartnerProbeSeconds));
var interval = TimeSpan.FromMilliseconds(Math.Max(50, _clusterOptions.BootstrapGuard.PartnerProbeIntervalMs));
var sw = Stopwatch.StartNew();
while (!ct.IsCancellationRequested)
{
if (await TryConnectAsync(host, port, ct).ConfigureAwait(false)) return true;
if (sw.Elapsed >= window) return false;
try { await Task.Delay(interval, ct).ConfigureAwait(false); }
catch (OperationCanceledException) { return false; }
}
return false;
}
private async Task<bool> TryConnectAsync(string host, int port, CancellationToken ct)
{
try
{
using var client = new TcpClient();
using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
cts.CancelAfter(Math.Max(100, _clusterOptions.BootstrapGuard.ProbeConnectTimeoutMs));
await client.ConnectAsync(host, port, cts.Token).ConfigureAwait(false);
return client.Connected;
}
catch (OperationCanceledException) when (ct.IsCancellationRequested)
{
throw; // host shutdown — propagate
}
catch
{
return false; // connect refused / timed out / DNS not resolvable yet — partner not up
}
}
}
+14
View File
@@ -394,6 +394,20 @@ try
builder.Services.AddSingleton<Akka.Actor.ActorSystem>(sp =>
sp.GetRequiredService<AkkaHostedService>().GetOrCreateActorSystem());
// Simultaneous-cold-start split-brain guard (Gitea #33, dark switch
// ScadaBridge:Cluster:BootstrapGuard:Enabled, default off). Registered AFTER the
// AkkaHostedService/ActorSystem bridge so its StartAsync resolves the system the main path
// built; when the guard is off it no-ops (Akka auto-joined from config seeds), so registering
// it unconditionally is safe. When on, the node started unjoined (empty HOCON seed list, see
// AkkaHostedService.BuildHocon) and this coordinator issues the single reachability-gated
// JoinSeedNodes. It resolves the ActorSystem lazily via GetOrCreateActorSystem so it never
// races startup or caches a null.
builder.Services.AddHostedService(sp => new ClusterBootstrapCoordinator(
() => sp.GetRequiredService<AkkaHostedService>().GetOrCreateActorSystem(),
sp.GetRequiredService<IOptions<ClusterOptions>>(),
sp.GetRequiredService<IOptions<NodeOptions>>(),
sp.GetRequiredService<ILogger<ClusterBootstrapCoordinator>>()));
// Register the production IActiveNodeGate implementation so
// standby-node gating is actually enforced (the InboundApiEndpointFilter
// consults IActiveNodeGate and defaults to "allow" when none is registered,
@@ -168,6 +168,19 @@ public static class SiteServiceRegistration
services.AddSingleton<Akka.Actor.ActorSystem>(sp =>
sp.GetRequiredService<AkkaHostedService>().GetOrCreateActorSystem());
// Simultaneous-cold-start split-brain guard (Gitea #33, dark switch
// ScadaBridge:Cluster:BootstrapGuard:Enabled, default off). This is the case the guard was
// built for — a shared power event that powers up both site VMs together. Registered AFTER
// the AkkaHostedService/ActorSystem bridge (mirrors the Central composition root in
// Program.cs); it no-ops when the guard is off, so registering it unconditionally is safe.
// When on, the node started unjoined (empty HOCON seed list) and this coordinator issues the
// single reachability-gated JoinSeedNodes, resolving the ActorSystem lazily.
services.AddHostedService(sp => new ClusterBootstrapCoordinator(
() => sp.GetRequiredService<AkkaHostedService>().GetOrCreateActorSystem(),
sp.GetRequiredService<IOptions<ClusterOptions>>(),
sp.GetRequiredService<IOptions<NodeOptions>>(),
sp.GetRequiredService<ILogger<ClusterBootstrapCoordinator>>()));
// Cluster node status provider for health reports
services.AddSingleton<IClusterNodeProvider>(sp =>
{