From 44ca7c8623f8a04cb20d60277964dc9cf9856ad7 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Sat, 15 Aug 2026 12:31:36 -0400 Subject: [PATCH] fix(dashboard): enforce/document the advised-set cap honestly; cover handle-0 and oversize-read paths --- docs/GatewayDashboardDesign.md | 13 ++ .../Dashboard/DashboardLiveDataService.cs | 33 ++++- .../DashboardLiveDataServiceTests.cs | 130 ++++++++++++++++-- 3 files changed, 163 insertions(+), 13 deletions(-) diff --git a/docs/GatewayDashboardDesign.md b/docs/GatewayDashboardDesign.md index d2327fa..77be9d4 100644 --- a/docs/GatewayDashboardDesign.md +++ b/docs/GatewayDashboardDesign.md @@ -406,6 +406,19 @@ room for each other. A failed unadvise does not fail the read: the tags are dropped from tracking anyway (they re-subscribe if read again), because the session-invalidation path already handles gateway/worker drift. +The cap is per-read, not absolute. A read may never evict a tag it is itself about +to return, so one read of more distinct tags than the cap leaves the set that +large; what the eviction pass guarantees is + +> after any read, the advise set holds at most `max(256, distinct tags in that read)` +> tags. + +The overshoot is not sticky: the next read that subscribes anything measures the +overflow against the oversized set and evicts the whole excess in one pass (a +300-tag set plus one new tag evicts 45 and lands back at 256). A read that +subscribes nothing new evicts nothing, but neither can it grow the set. A browse +page requests far fewer tags than the cap, so in practice the set settles at 256. + The Alarms page does **not** use the dashboard session: alarm data comes from the gateway's always-on central monitor. `QueryAlarmsAsync` reads `IGatewayAlarmService.CurrentAlarms` — the monitor's in-process cache — so the diff --git a/src/ZB.MOM.WW.MxGateway.Server/Dashboard/DashboardLiveDataService.cs b/src/ZB.MOM.WW.MxGateway.Server/Dashboard/DashboardLiveDataService.cs index d28a44c..ec2aa2e 100644 --- a/src/ZB.MOM.WW.MxGateway.Server/Dashboard/DashboardLiveDataService.cs +++ b/src/ZB.MOM.WW.MxGateway.Server/Dashboard/DashboardLiveDataService.cs @@ -19,6 +19,16 @@ public sealed class DashboardLiveDataService : IDashboardLiveDataService, IAsync // One browse page of tags plus headroom. Bounds the standing advise load the // single dashboard worker carries — and the event churn that advise set feeds — // however much of a galaxy an operator browses through in one sitting. + // + // The bound is per-read, not absolute: a read may never evict a tag it is itself + // about to return, so a single read of more distinct tags than the cap leaves the + // set that large. The invariant EvictForAsync actually maintains is + // + // |advise set| after a read <= max(MaxSubscribedTags, distinct tags in that read) + // + // and any overshoot is squeezed back out by the next read that subscribes a tag + // (see EvictForAsync). A browse page requests far fewer tags than the cap, so in + // practice the set settles at MaxSubscribedTags. private const int MaxSubscribedTags = 256; private static readonly TimeSpan ReadTimeout = TimeSpan.FromSeconds(5); @@ -125,6 +135,11 @@ public sealed class DashboardLiveDataService : IDashboardLiveDataService, IAsync // order). `justReadCount` is how many distinct tags of this read were already // advised — they now occupy the front of the list and must never be evicted to // make room for the same read's new tags. Callers must hold _gate. + // + // Every tag of one read is equally recently read; the recency list needs a total + // order anyway, so the whole service uses one tie-break: later in the request wins. + // Promoting in request order gives that here, and TrackSubscribed inserts new tags + // the same way. private string[] TouchAndCollectNewTags(IReadOnlyCollection tagAddresses, out int justReadCount) { int touched = 0; @@ -161,6 +176,18 @@ public sealed class DashboardLiveDataService : IDashboardLiveDataService, IAsync // one batch. A failed unadvise must not fail the read: the tags are dropped // from tracking regardless, and the session-invalidation path already handles // gateway/worker drift. Callers must hold _gate. + // + // Eviction stops at the tags this read just touched (`justReadCount`), so a read + // whose own distinct tags outnumber the cap ends over it — see MaxSubscribedTags + // for the exact invariant. That overshoot is not sticky: the next read that + // subscribes anything computes `overflow` against the oversized set and evicts the + // whole excess in one pass (a 300-tag set plus one new tag evicts 45 and lands + // back at the cap). A read that subscribes nothing new evicts nothing, but it also + // cannot grow the set. + // + // Cancellation mid-eviction follows this file's policy: OperationCanceledException + // is deliberately not caught here or in ReadAsync, so it propagates with the tags + // already dropped from tracking — the same end state as a failed unadvise. private async Task EvictForAsync( GatewaySession session, int serverHandle, @@ -227,10 +254,10 @@ public sealed class DashboardLiveDataService : IDashboardLiveDataService, IAsync } } - // Inserted back-to-front so the read's first tag ends up most recent. - for (int i = tagAddresses.Count - 1; i >= 0; i--) + // Request order, so the read's last tag ends up most recent — the same + // tie-break TouchAndCollectNewTags applies to the tags it promotes. + foreach (string tag in tagAddresses) { - string tag = tagAddresses[i]; handles.TryGetValue(tag, out int itemHandle); _subscribed[tag] = _recency.AddFirst(new SubscribedTag(tag, itemHandle)); } diff --git a/src/ZB.MOM.WW.MxGateway.Tests/Gateway/Dashboard/DashboardLiveDataServiceTests.cs b/src/ZB.MOM.WW.MxGateway.Tests/Gateway/Dashboard/DashboardLiveDataServiceTests.cs index dcb12ca..3059458 100644 --- a/src/ZB.MOM.WW.MxGateway.Tests/Gateway/Dashboard/DashboardLiveDataServiceTests.cs +++ b/src/ZB.MOM.WW.MxGateway.Tests/Gateway/Dashboard/DashboardLiveDataServiceTests.cs @@ -46,12 +46,14 @@ public sealed class DashboardLiveDataServiceTests await using FakeSessionManager sessionManager = new(worker); await using DashboardLiveDataService service = CreateService(sessionManager); + // Within one read the last tag counts as most recently read, so filler[0] is + // the tail of the recency list. string[] filler = CreateTagAddresses(MaxSubscribedTags); await service.ReadAsync(filler, CancellationToken.None); - // Re-read the tag that is currently least recent, making it most recent: the - // tag read before it is now the eviction candidate. - await service.ReadAsync([filler[^1]], CancellationToken.None); + // Re-read the tail, making it most recent: the tag read before it is now the + // eviction candidate. + await service.ReadAsync([filler[0]], CancellationToken.None); Assert.Equal(MaxSubscribedTags, worker.SubscribedTags.Count); DashboardLiveReadResult overflow = await service.ReadAsync( @@ -59,16 +61,104 @@ public sealed class DashboardLiveDataServiceTests CancellationToken.None); Assert.Null(overflow.Error); - Assert.Equal([worker.HandleFor(filler[^2])], worker.UnsubscribedHandles); + Assert.Equal([worker.HandleFor(filler[1])], worker.UnsubscribedHandles); Assert.Equal("Overflow.PV", worker.SubscribedTags[^1]); Assert.Equal(MaxSubscribedTags + 1, worker.SubscribedTags.Count); // The evicted tag is no longer tracked and re-subscribes; the re-read one does not. - await service.ReadAsync([filler[^1], filler[^2]], CancellationToken.None); - Assert.Equal(filler[^2], worker.SubscribedTags[^1]); + await service.ReadAsync([filler[0], filler[1]], CancellationToken.None); + Assert.Equal(filler[1], worker.SubscribedTags[^1]); Assert.Equal(MaxSubscribedTags + 2, worker.SubscribedTags.Count); } + /// + /// Verifies the cap is per-read, not absolute: a single read of more distinct tags + /// than the cap keeps them all (a read never evicts a tag it is about to return), + /// and the next read that subscribes anything squeezes the overshoot back out. + /// + [Fact] + public async Task ReadAsync_WithMoreDistinctTagsThanCap_KeepsThemAllThenSelfCorrects() + { + RecordingWorkerClient worker = new(); + await using FakeSessionManager sessionManager = new(worker); + await using DashboardLiveDataService service = CreateService(sessionManager); + + string[] oversize = CreateTagAddresses(300); + DashboardLiveReadResult oversizeResult = await service.ReadAsync(oversize, CancellationToken.None); + + Assert.Null(oversizeResult.Error); + Assert.Equal(300, oversizeResult.Values.Count); + Assert.Equal(300, worker.SubscribedTags.Count); + Assert.Empty(worker.UnsubscribedHandles); + + // 300 + 1 - 256 = 45 evicted in one pass, landing the set back on the cap. + await service.ReadAsync(["Overflow.PV"], CancellationToken.None); + Assert.Equal(oversize[..45].Select(worker.HandleFor), worker.UnsubscribedHandles); + + // Exactly at the cap now: one more new tag evicts exactly one. + worker.UnsubscribedHandles.Clear(); + await service.ReadAsync(["Overflow2.PV"], CancellationToken.None); + Assert.Equal([worker.HandleFor(oversize[45])], worker.UnsubscribedHandles); + } + + /// + /// Verifies tags read in the same call are never evicted for each other: a read that + /// touches nearly the whole advise set evicts only the untouched remainder, ends over + /// the cap, and the following read trims it back. + /// + [Fact] + public async Task ReadAsync_WhenTouchedTagsFillTheCap_EvictsOnlyUntouchedTags() + { + RecordingWorkerClient worker = new(); + await using FakeSessionManager sessionManager = new(worker); + await using DashboardLiveDataService service = CreateService(sessionManager); + + string[] filler = CreateTagAddresses(MaxSubscribedTags); + await service.ReadAsync(filler, CancellationToken.None); + + // 250 already-advised tags + 10 new ones: only the 6 untouched tags are + // evictable, so the set ends at 260. + string[] fresh = CreateTagAddresses(10, "Fresh"); + await service.ReadAsync([.. filler[..250], .. fresh], CancellationToken.None); + + Assert.Equal(filler[250..].Select(worker.HandleFor), worker.UnsubscribedHandles); + + // 260 + 1 - 256 = 5 evicted on the next read that subscribes anything. + worker.UnsubscribedHandles.Clear(); + await service.ReadAsync(["Overflow.PV"], CancellationToken.None); + Assert.Equal(5, worker.UnsubscribedHandles.Count); + } + + /// + /// Verifies a tag the worker failed to advise still occupies a slot but is evicted + /// without any unsubscribe command — there is no item handle to unadvise. + /// + [Fact] + public async Task ReadAsync_WhenAdviseFailed_EvictsTagWithoutUnsubscribing() + { + RecordingWorkerClient worker = new(); + worker.FailSubscribeFor.Add("Bad.PV"); + await using FakeSessionManager sessionManager = new(worker); + await using DashboardLiveDataService service = CreateService(sessionManager); + + // Bad.PV is read first, so it is the least recently read of the batch and the + // first tag evicted. + string[] filler = CreateTagAddresses(MaxSubscribedTags - 1); + await service.ReadAsync(["Bad.PV", .. filler], CancellationToken.None); + + DashboardLiveReadResult overflow = await service.ReadAsync( + ["Overflow.PV"], + CancellationToken.None); + + Assert.Null(overflow.Error); + Assert.Empty(worker.UnsubscribedHandles); + Assert.Equal(0, worker.UnsubscribeCommandCount); + + // It was dropped from tracking all the same, so reading it again re-advises it. + await service.ReadAsync(["Bad.PV"], CancellationToken.None); + Assert.Equal("Bad.PV", worker.SubscribedTags[^1]); + } + /// /// Verifies a failed unadvise of an evicted tag does not fail the read, and the /// evicted tag is dropped from tracking anyway. @@ -92,8 +182,8 @@ public sealed class DashboardLiveDataServiceTests Assert.Equal(1, sessionManager.OpenCount); // The evicted tag was dropped from tracking despite the failed unadvise. - await service.ReadAsync([filler[^1]], CancellationToken.None); - Assert.Equal(filler[^1], worker.SubscribedTags[^1]); + await service.ReadAsync([filler[0]], CancellationToken.None); + Assert.Equal(filler[0], worker.SubscribedTags[^1]); } private static DashboardLiveDataService CreateService(ISessionManager sessionManager) @@ -104,12 +194,12 @@ public sealed class DashboardLiveDataServiceTests NullLogger.Instance); } - private static string[] CreateTagAddresses(int count) + private static string[] CreateTagAddresses(int count, string prefix = "Tank") { string[] addresses = new string[count]; for (int i = 0; i < count; i++) { - addresses[i] = $"Tank_{i:D4}.PV"; + addresses[i] = $"{prefix}_{i:D4}.PV"; } return addresses; @@ -227,6 +317,12 @@ public sealed class DashboardLiveDataServiceTests /// Gets the item handles the dashboard unsubscribed, in order. public List UnsubscribedHandles { get; } = []; + /// Gets the number of unsubscribe commands the dashboard sent. + public int UnsubscribeCommandCount { get; private set; } + + /// Gets the tag addresses the worker refuses to advise. + public HashSet FailSubscribeFor { get; } = new(StringComparer.OrdinalIgnoreCase); + /// Gets or sets a value indicating whether unsubscribe commands throw. public bool FailUnsubscribe { get; set; } @@ -262,6 +358,7 @@ public sealed class DashboardLiveDataServiceTests reply.SubscribeBulk = Subscribe(mxCommand.SubscribeBulk.TagAddresses); break; case MxCommandKind.UnsubscribeBulk: + UnsubscribeCommandCount++; if (FailUnsubscribe) { throw new InvalidOperationException("Simulated worker unsubscribe failure."); @@ -305,6 +402,19 @@ public sealed class DashboardLiveDataServiceTests foreach (string tagAddress in tagAddresses) { SubscribedTags.Add(tagAddress); + if (FailSubscribeFor.Contains(tagAddress)) + { + subscribeReply.Results.Add(new SubscribeResult + { + ServerHandle = RegisteredServerHandle, + TagAddress = tagAddress, + ItemHandle = 0, + WasSuccessful = false, + ErrorMessage = "Simulated advise failure.", + }); + continue; + } + if (!_itemHandles.TryGetValue(tagAddress, out int itemHandle)) { itemHandle = _nextItemHandle++;