fix(health): evict deleted sites from the aggregator on the periodic site refresh
This commit is contained in:
@@ -499,7 +499,12 @@ public class CentralCommunicationActor : ReceiveActor
|
|||||||
var frozen = contacts.ToDictionary(
|
var frozen = contacts.ToDictionary(
|
||||||
kvp => kvp.Key,
|
kvp => kvp.Key,
|
||||||
kvp => (IReadOnlyList<string>)kvp.Value.AsReadOnly());
|
kvp => (IReadOnlyList<string>)kvp.Value.AsReadOnly());
|
||||||
return new SiteAddressCacheLoaded(frozen);
|
|
||||||
|
// Carry the full configured-site identifier set (not just the
|
||||||
|
// address-bearing subset in `frozen`) so the aggregator prunes only
|
||||||
|
// genuinely-deleted sites and never an addressless-but-configured one.
|
||||||
|
var knownSiteIds = sites.Select(s => s.SiteIdentifier).ToList();
|
||||||
|
return new SiteAddressCacheLoaded(frozen, knownSiteIds);
|
||||||
}).PipeTo(self);
|
}).PipeTo(self);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -570,6 +575,12 @@ public class CentralCommunicationActor : ReceiveActor
|
|||||||
}
|
}
|
||||||
|
|
||||||
_log.Info("Site ClusterClient cache refreshed with {0} site(s)", _siteClients.Count);
|
_log.Info("Site ClusterClient cache refreshed with {0} site(s)", _siteClients.Count);
|
||||||
|
|
||||||
|
// Self-healing eviction: a site deleted from configuration would otherwise
|
||||||
|
// linger in the aggregator as a permanently-offline tile (and live KPI
|
||||||
|
// sample source) forever. Prune on every refresh so it disappears within
|
||||||
|
// one refresh interval without needing a dedicated deletion event.
|
||||||
|
_serviceProvider.GetService<ICentralHealthAggregator>()?.PruneUnknownSites(msg.KnownSiteIds);
|
||||||
}
|
}
|
||||||
|
|
||||||
// TrackMessageForCleanup removed — the dicts it fed
|
// TrackMessageForCleanup removed — the dicts it fed
|
||||||
@@ -650,7 +661,9 @@ public record RefreshSiteAddresses;
|
|||||||
/// discipline. The producer wraps the constructed buckets with
|
/// discipline. The producer wraps the constructed buckets with
|
||||||
/// <c>List<T>.AsReadOnly()</c> before piping to Self.
|
/// <c>List<T>.AsReadOnly()</c> before piping to Self.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
internal record SiteAddressCacheLoaded(IReadOnlyDictionary<string, IReadOnlyList<string>> SiteContacts);
|
internal record SiteAddressCacheLoaded(
|
||||||
|
IReadOnlyDictionary<string, IReadOnlyList<string>> SiteContacts,
|
||||||
|
IReadOnlyCollection<string> KnownSiteIds);
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Notification sent to debug view subscribers when the stream is terminated
|
/// Notification sent to debug view subscribers when the stream is terminated
|
||||||
|
|||||||
@@ -196,6 +196,20 @@ public class CentralHealthAggregator : BackgroundService, ICentralHealthAggregat
|
|||||||
return state;
|
return state;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public void PruneUnknownSites(IReadOnlyCollection<string> knownSiteIds)
|
||||||
|
{
|
||||||
|
var known = new HashSet<string>(knownSiteIds, StringComparer.Ordinal);
|
||||||
|
foreach (var siteId in _siteStates.Keys)
|
||||||
|
{
|
||||||
|
if (siteId == CentralHealthReportLoop.CentralSiteId || known.Contains(siteId))
|
||||||
|
continue;
|
||||||
|
if (_siteStates.TryRemove(siteId, out _))
|
||||||
|
_logger.LogInformation(
|
||||||
|
"Site {SiteId} evicted from health aggregator (no longer configured)", siteId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -39,4 +39,16 @@ public interface ICentralHealthAggregator
|
|||||||
/// <param name="siteId">The string identifier of the site to look up.</param>
|
/// <param name="siteId">The string identifier of the site to look up.</param>
|
||||||
/// <returns>The <see cref="SiteHealthState"/> for the site, or null if not yet seen.</returns>
|
/// <returns>The <see cref="SiteHealthState"/> for the site, or null if not yet seen.</returns>
|
||||||
SiteHealthState? GetSiteState(string siteId);
|
SiteHealthState? GetSiteState(string siteId);
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Removes tracked sites not in <paramref name="knownSiteIds"/> (i.e. deleted
|
||||||
|
/// from configuration). The synthetic central id is always retained.
|
||||||
|
/// Self-healing: called from the CentralCommunicationActor's periodic site
|
||||||
|
/// refresh, so a deleted site disappears within one refresh interval without
|
||||||
|
/// needing a dedicated deletion event. Without this the in-memory state only
|
||||||
|
/// ever grew, leaving a deleted site as a permanently-offline dashboard tile
|
||||||
|
/// (and a live KPI sample source) forever.
|
||||||
|
/// </summary>
|
||||||
|
/// <param name="knownSiteIds">The identifiers of all currently-configured sites.</param>
|
||||||
|
void PruneUnknownSites(IReadOnlyCollection<string> knownSiteIds);
|
||||||
}
|
}
|
||||||
|
|||||||
+32
-1
@@ -10,6 +10,7 @@ using ZB.MOM.WW.ScadaBridge.Commons.Entities.Sites;
|
|||||||
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories;
|
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories;
|
||||||
using ZB.MOM.WW.ScadaBridge.Communication;
|
using ZB.MOM.WW.ScadaBridge.Communication;
|
||||||
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
|
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
|
||||||
|
using ZB.MOM.WW.ScadaBridge.HealthMonitoring;
|
||||||
|
|
||||||
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
|
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
|
||||||
|
|
||||||
@@ -37,7 +38,37 @@ public class CentralCommunicationActorClientLifecycleTests : TestKit
|
|||||||
|
|
||||||
private static SiteAddressCacheLoaded Load(string siteId, params string[] addrs) =>
|
private static SiteAddressCacheLoaded Load(string siteId, params string[] addrs) =>
|
||||||
new(new Dictionary<string, IReadOnlyList<string>>
|
new(new Dictionary<string, IReadOnlyList<string>>
|
||||||
{ [siteId] = addrs.ToList().AsReadOnly() });
|
{ [siteId] = addrs.ToList().AsReadOnly() },
|
||||||
|
new[] { siteId });
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public void PeriodicRefresh_PrunesDeletedSites_FromHealthAggregator()
|
||||||
|
{
|
||||||
|
// PLAN-01 Task 18: the periodic site refresh feeds the currently-configured
|
||||||
|
// site ids to the aggregator so a deleted site is evicted within one
|
||||||
|
// refresh interval. Repo returns only "site-a".
|
||||||
|
var siteRepo = Substitute.For<ISiteRepository>();
|
||||||
|
siteRepo.GetAllSitesAsync(Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new List<Site> { new("Site A", "site-a") });
|
||||||
|
var aggregator = Substitute.For<ICentralHealthAggregator>();
|
||||||
|
var services = new ServiceCollection();
|
||||||
|
services.AddScoped(_ => siteRepo);
|
||||||
|
services.AddSingleton(aggregator);
|
||||||
|
var provider = services.BuildServiceProvider();
|
||||||
|
|
||||||
|
var actor = Sys.ActorOf(Props.Create(() => new CentralCommunicationActor(
|
||||||
|
provider, new DefaultSiteClientFactory(), null)));
|
||||||
|
|
||||||
|
// Trigger the refresh (also fires at PreStart, but drive it explicitly so
|
||||||
|
// the assertion is deterministic). The load runs on a detached task and
|
||||||
|
// pipes SiteAddressCacheLoaded back to the actor, which then prunes.
|
||||||
|
actor.Tell(new RefreshSiteAddresses());
|
||||||
|
|
||||||
|
AwaitAssert(
|
||||||
|
() => aggregator.Received().PruneUnknownSites(
|
||||||
|
Arg.Is<IReadOnlyCollection<string>>(ids => ids.Contains("site-a"))),
|
||||||
|
TimeSpan.FromSeconds(3));
|
||||||
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
public void AddressEdit_RecreatesClient_WithoutRestartLoop()
|
public void AddressEdit_RecreatesClient_WithoutRestartLoop()
|
||||||
|
|||||||
@@ -473,6 +473,27 @@ public class CentralHealthAggregatorTests
|
|||||||
Assert.False(aggregator.GetSiteState("site-a")!.IsMetricsStale);
|
Assert.False(aggregator.GetSiteState("site-a")!.IsMetricsStale);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Review 01 [Low]: <c>_siteStates</c> only ever grew — a site deleted from
|
||||||
|
/// configuration remained a permanently-offline dashboard tile (and a KPI
|
||||||
|
/// sample source) forever. Pruning against the known-site set evicts it; the
|
||||||
|
/// synthetic central id is always retained.
|
||||||
|
/// </summary>
|
||||||
|
[Fact]
|
||||||
|
public void PruneUnknownSites_RemovesDeletedSites_KeepsKnownAndCentral()
|
||||||
|
{
|
||||||
|
var (aggregator, _) = NewAggregator();
|
||||||
|
aggregator.ProcessReport(MakeReport("site-a", 1));
|
||||||
|
aggregator.ProcessReport(MakeReport("site-b", 1));
|
||||||
|
aggregator.ProcessReport(MakeReport(CentralHealthReportLoop.CentralSiteId, 1));
|
||||||
|
|
||||||
|
aggregator.PruneUnknownSites(new[] { "site-a" });
|
||||||
|
|
||||||
|
Assert.NotNull(aggregator.GetSiteState("site-a"));
|
||||||
|
Assert.Null(aggregator.GetSiteState("site-b"));
|
||||||
|
Assert.NotNull(aggregator.GetSiteState(CentralHealthReportLoop.CentralSiteId)); // synthetic id always kept
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Review 01 underdeveloped #7: "when did the site drop" was unanswerable.
|
/// Review 01 underdeveloped #7: "when did the site drop" was unanswerable.
|
||||||
/// Every online↔offline flip stamps <see cref="SiteHealthState.LastStatusChangeAt"/>.
|
/// Every online↔offline flip stamps <see cref="SiteHealthState.LastStatusChangeAt"/>.
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ public class CentralHealthReportLoopTests
|
|||||||
public IReadOnlyDictionary<string, SiteHealthState> GetAllSiteStates() =>
|
public IReadOnlyDictionary<string, SiteHealthState> GetAllSiteStates() =>
|
||||||
new Dictionary<string, SiteHealthState>();
|
new Dictionary<string, SiteHealthState>();
|
||||||
public SiteHealthState? GetSiteState(string siteId) => null;
|
public SiteHealthState? GetSiteState(string siteId) => null;
|
||||||
|
public void PruneUnknownSites(IReadOnlyCollection<string> knownSiteIds) { }
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -256,6 +257,7 @@ public class CentralHealthReportLoopTests
|
|||||||
public IReadOnlyDictionary<string, SiteHealthState> GetAllSiteStates() =>
|
public IReadOnlyDictionary<string, SiteHealthState> GetAllSiteStates() =>
|
||||||
new Dictionary<string, SiteHealthState>();
|
new Dictionary<string, SiteHealthState>();
|
||||||
public SiteHealthState? GetSiteState(string siteId) => null;
|
public SiteHealthState? GetSiteState(string siteId) => null;
|
||||||
|
public void PruneUnknownSites(IReadOnlyCollection<string> knownSiteIds) { }
|
||||||
}
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
|
|||||||
+3
@@ -166,5 +166,8 @@ public class SiteHealthKpiSampleSourceTests
|
|||||||
|
|
||||||
public void MarkHeartbeat(string siteId, DateTimeOffset receivedAt) =>
|
public void MarkHeartbeat(string siteId, DateTimeOffset receivedAt) =>
|
||||||
throw new NotSupportedException();
|
throw new NotSupportedException();
|
||||||
|
|
||||||
|
public void PruneUnknownSites(IReadOnlyCollection<string> knownSiteIds) =>
|
||||||
|
throw new NotSupportedException();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user