feat(site-runtime): wire ConfigFetchRetryCount into the standby replicated-config fetch (UA2)
This commit is contained in:
@@ -769,7 +769,7 @@ akka {{
|
|||||||
var replicationActor = _actorSystem!.ActorOf(
|
var replicationActor = _actorSystem!.ActorOf(
|
||||||
Props.Create(() => new SiteReplicationActor(
|
Props.Create(() => new SiteReplicationActor(
|
||||||
storage, sfStorage, replicationService, siteRole, replicationLogger,
|
storage, sfStorage, replicationService, siteRole, replicationLogger,
|
||||||
deploymentConfigFetcher)),
|
deploymentConfigFetcher, null, siteRuntimeOptionsValue, null)),
|
||||||
"site-replication");
|
"site-replication");
|
||||||
|
|
||||||
// Wire S&F replication handler to forward operations via the replication actor
|
// Wire S&F replication handler to forward operations via the replication actor
|
||||||
|
|||||||
@@ -28,6 +28,8 @@ public class SiteReplicationActor : ReceiveActor
|
|||||||
private readonly ILogger<SiteReplicationActor> _logger;
|
private readonly ILogger<SiteReplicationActor> _logger;
|
||||||
private readonly Cluster _cluster;
|
private readonly Cluster _cluster;
|
||||||
private readonly Func<bool> _isActive;
|
private readonly Func<bool> _isActive;
|
||||||
|
private readonly int _configFetchRetryCount;
|
||||||
|
private readonly TimeSpan _configFetchRetryDelay;
|
||||||
private Address? _peerAddress;
|
private Address? _peerAddress;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -64,7 +66,9 @@ public class SiteReplicationActor : ReceiveActor
|
|||||||
string siteRole,
|
string siteRole,
|
||||||
ILogger<SiteReplicationActor> logger,
|
ILogger<SiteReplicationActor> logger,
|
||||||
IDeploymentConfigFetcher? configFetcher = null,
|
IDeploymentConfigFetcher? configFetcher = null,
|
||||||
Func<bool>? isActiveOverride = null)
|
Func<bool>? isActiveOverride = null,
|
||||||
|
SiteRuntimeOptions? options = null,
|
||||||
|
TimeSpan? configFetchRetryDelay = null)
|
||||||
{
|
{
|
||||||
_storage = storage;
|
_storage = storage;
|
||||||
_sfStorage = sfStorage;
|
_sfStorage = sfStorage;
|
||||||
@@ -74,6 +78,11 @@ public class SiteReplicationActor : ReceiveActor
|
|||||||
_logger = logger;
|
_logger = logger;
|
||||||
_cluster = Cluster.Get(Context.System);
|
_cluster = Cluster.Get(Context.System);
|
||||||
_isActive = isActiveOverride ?? DefaultIsActive;
|
_isActive = isActiveOverride ?? DefaultIsActive;
|
||||||
|
// UA2: bound the standby's replicated-config fetch retries. At least one
|
||||||
|
// attempt always runs; the fixed inter-attempt delay is a test seam
|
||||||
|
// (production default 2 s).
|
||||||
|
_configFetchRetryCount = Math.Max(1, options?.ConfigFetchRetryCount ?? 1);
|
||||||
|
_configFetchRetryDelay = configFetchRetryDelay ?? TimeSpan.FromSeconds(2);
|
||||||
|
|
||||||
// Cluster member events
|
// Cluster member events
|
||||||
Receive<ClusterEvent.MemberUp>(HandleMemberUp);
|
Receive<ClusterEvent.MemberUp>(HandleMemberUp);
|
||||||
@@ -254,43 +263,53 @@ public class SiteReplicationActor : ReceiveActor
|
|||||||
|
|
||||||
// Notify-and-fetch: the peer sent only the id, so the standby fetches the config
|
// Notify-and-fetch: the peer sent only the id, so the standby fetches the config
|
||||||
// itself (off-thread; best-effort fire-and-forget, matching the no-ack replication
|
// itself (off-thread; best-effort fire-and-forget, matching the no-ack replication
|
||||||
// model). The guarded write only overwrites a strictly-older local row. A single
|
// model). The guarded write only overwrites a strictly-older local row. The fetch
|
||||||
// fetch attempt — reconciliation is the durable backstop for a lost fetch.
|
// is retried up to ConfigFetchRetryCount times with a fixed delay (UA2) — a transient
|
||||||
_configFetcher.FetchAsync(msg.CentralFetchBaseUrl, msg.DeploymentId, msg.FetchToken, CancellationToken.None)
|
// central hiccup no longer defers to the slower reconciliation backstop, which still
|
||||||
.ContinueWith(async t =>
|
// covers a total failure after the last attempt.
|
||||||
|
_ = FetchWithRetryAsync();
|
||||||
|
return;
|
||||||
|
|
||||||
|
async Task FetchWithRetryAsync()
|
||||||
|
{
|
||||||
|
for (var attempt = 1; attempt <= _configFetchRetryCount; attempt++)
|
||||||
{
|
{
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
if (t.IsCompletedSuccessfully)
|
// Non-null: the outer method returns early when _configFetcher is null.
|
||||||
{
|
var json = await _configFetcher!.FetchAsync(
|
||||||
await _storage.StoreDeployedConfigIfNewerAsync(
|
msg.CentralFetchBaseUrl, msg.DeploymentId, msg.FetchToken, CancellationToken.None);
|
||||||
msg.InstanceName, t.Result, msg.DeploymentId, msg.RevisionHash, msg.IsEnabled);
|
await _storage.StoreDeployedConfigIfNewerAsync(
|
||||||
return;
|
msg.InstanceName, json, msg.DeploymentId, msg.RevisionHash, msg.IsEnabled);
|
||||||
}
|
return;
|
||||||
|
|
||||||
var ex = t.Exception?.GetBaseException();
|
|
||||||
if (ex is DeploymentConfigFetchException { IsSuperseded: true })
|
|
||||||
_logger.LogInformation(
|
|
||||||
"Skip replicated config for {Instance}: superseded/expired (a newer deploy will replicate)",
|
|
||||||
msg.InstanceName);
|
|
||||||
else if (t.IsCanceled)
|
|
||||||
_logger.LogWarning(
|
|
||||||
"Replicated config fetch cancelled for {Instance} (deployment {DeploymentId})",
|
|
||||||
msg.InstanceName, msg.DeploymentId);
|
|
||||||
else
|
|
||||||
_logger.LogError(ex,
|
|
||||||
"Replicated config fetch failed for {Instance} (deployment {DeploymentId})",
|
|
||||||
msg.InstanceName, msg.DeploymentId);
|
|
||||||
}
|
}
|
||||||
catch (Exception writeEx)
|
catch (DeploymentConfigFetchException fex) when (fex.IsSuperseded)
|
||||||
{
|
{
|
||||||
// Guarded-write failure is best-effort; observe + log so nothing faults silently.
|
// A superseded/expired fetch never heals by retrying — a newer deploy
|
||||||
_logger.LogError(writeEx,
|
// will replicate its own id.
|
||||||
"Failed to write replicated config for {Instance} (deployment {DeploymentId})",
|
_logger.LogInformation(
|
||||||
msg.InstanceName, msg.DeploymentId);
|
"Skip replicated config for {Instance}: superseded/expired (a newer deploy will replicate)",
|
||||||
|
msg.InstanceName);
|
||||||
|
return;
|
||||||
}
|
}
|
||||||
})
|
catch (Exception ex)
|
||||||
.Unwrap();
|
{
|
||||||
|
if (attempt < _configFetchRetryCount)
|
||||||
|
{
|
||||||
|
_logger.LogWarning(ex,
|
||||||
|
"Replicated config fetch attempt {Attempt}/{Max} failed for {Instance} (deployment {DeploymentId}) — retrying in {Delay}",
|
||||||
|
attempt, _configFetchRetryCount, msg.InstanceName, msg.DeploymentId, _configFetchRetryDelay);
|
||||||
|
await Task.Delay(_configFetchRetryDelay);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
_logger.LogError(ex,
|
||||||
|
"Replicated config fetch failed after {Attempts} attempt(s) for {Instance} (deployment {DeploymentId}) — reconciliation will backstop",
|
||||||
|
_configFetchRetryCount, msg.InstanceName, msg.DeploymentId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void HandleApplyConfigRemove(ApplyConfigRemove msg)
|
private void HandleApplyConfigRemove(ApplyConfigRemove msg)
|
||||||
|
|||||||
@@ -61,8 +61,9 @@ public class SiteRuntimeOptions
|
|||||||
public int ConfigFetchTimeoutSeconds { get; set; } = 30;
|
public int ConfigFetchTimeoutSeconds { get; set; } = 30;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Bounded retry count for the standby's best-effort replicated-config fetch.
|
/// Bounded attempt count (including the first) for the standby's replicated-config
|
||||||
/// Reserved — consumed by the standby replication fetch in a later task; not yet wired.
|
/// fetch; a 2 s fixed delay separates attempts and superseded fetches never retry.
|
||||||
|
/// Consumed by <c>SiteReplicationActor.HandleApplyConfigDeploy</c> (UA2). Default: 3.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public int ConfigFetchRetryCount { get; set; } = 3;
|
public int ConfigFetchRetryCount { get; set; } = 3;
|
||||||
|
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ using Akka.Actor;
|
|||||||
using Akka.TestKit.Xunit2;
|
using Akka.TestKit.Xunit2;
|
||||||
using Microsoft.Extensions.Logging;
|
using Microsoft.Extensions.Logging;
|
||||||
using Microsoft.Extensions.Logging.Abstractions;
|
using Microsoft.Extensions.Logging.Abstractions;
|
||||||
|
using ZB.MOM.WW.ScadaBridge.SiteRuntime;
|
||||||
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
|
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
|
||||||
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Deployment;
|
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Deployment;
|
||||||
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Messages;
|
using ZB.MOM.WW.ScadaBridge.SiteRuntime.Messages;
|
||||||
@@ -79,6 +80,39 @@ akka {
|
|||||||
_storage, _sfStorage, _replicationService, SiteRole,
|
_storage, _sfStorage, _replicationService, SiteRole,
|
||||||
NullLogger<SiteReplicationActor>.Instance, fetcher)));
|
NullLogger<SiteReplicationActor>.Instance, fetcher)));
|
||||||
|
|
||||||
|
private IActorRef CreateReplicationActor(
|
||||||
|
IDeploymentConfigFetcher fetcher, SiteRuntimeOptions options, TimeSpan retryDelay) =>
|
||||||
|
ActorOf(Props.Create(() => new SiteReplicationActor(
|
||||||
|
_storage, _sfStorage, _replicationService, SiteRole,
|
||||||
|
NullLogger<SiteReplicationActor>.Instance, fetcher, null, options, retryDelay)));
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task ReplicatedFetch_RetriesUpToConfigFetchRetryCount()
|
||||||
|
{
|
||||||
|
// The first two fetches fail transiently; the third succeeds. With
|
||||||
|
// ConfigFetchRetryCount = 3 the standby must retry to the third attempt and
|
||||||
|
// then guarded-write the fetched config (a short retry delay keeps the test fast).
|
||||||
|
var attempts = 0;
|
||||||
|
var fetcher = new FakeConfigFetcher(_ =>
|
||||||
|
Interlocked.Increment(ref attempts) < 3
|
||||||
|
? Task.FromException<string>(new InvalidOperationException("central hiccup"))
|
||||||
|
: Task.FromResult("{\"instanceUniqueName\":\"RetryPump\"}"));
|
||||||
|
var actor = CreateReplicationActor(
|
||||||
|
fetcher, new SiteRuntimeOptions { ConfigFetchRetryCount = 3 },
|
||||||
|
TimeSpan.FromMilliseconds(50));
|
||||||
|
|
||||||
|
actor.Tell(new ApplyConfigDeploy(
|
||||||
|
"RetryPump", "dep-r1", "sha256:r1", true,
|
||||||
|
"http://central:9000", "tok-r1"));
|
||||||
|
|
||||||
|
await AwaitAssertAsync(async () =>
|
||||||
|
{
|
||||||
|
Assert.Equal(3, Volatile.Read(ref attempts));
|
||||||
|
var configs = await _storage.GetAllDeployedConfigsAsync();
|
||||||
|
Assert.Single(configs, c => c.InstanceUniqueName == "RetryPump");
|
||||||
|
}, TimeSpan.FromSeconds(10));
|
||||||
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
public async Task ApplyConfigDeploy_StandbyFetchesConfigAndGuardedWrites()
|
public async Task ApplyConfigDeploy_StandbyFetchesConfigAndGuardedWrites()
|
||||||
{
|
{
|
||||||
|
|||||||
Reference in New Issue
Block a user