Merge branch 'fix/gwc-26-27-alarm-attach'
ci / nightly-windev (push) Has been skipped
ci / windows-x86 (push) Successful in 1m17s
ci / java (push) Successful in 2m4s
ci / portable (push) Successful in 7m8s

# Conflicts:
#	archreview/2026-07-12/remediation/00-tracking.md
#	archreview/2026-07-12/remediation/10-gateway-core.md
This commit is contained in:
Joseph Doherty
2026-08-07 06:02:17 -04:00
19 changed files with 738 additions and 130 deletions
@@ -61,8 +61,8 @@ Full design + implementation for each row lives in the linked domain doc under i
|---|---|:-:|:-:|---|---|---|
| GWC-24 | Medium | P1 | M | GWC-21 (coord, old tracker) | Done | Unbounded event staging channel: sustained slow drain grows memory silently and invisibly |
| GWC-25 | Medium | P0 | S | CLI-35/36 (coord) | Done | Empty-ring ReplayGap sentinel carries `oldest_available_sequence = 0` |
| GWC-26 | Low | P2 | M | GWC-27 | Not started | Alarm monitor attaches its subscriber after SubscribeAlarms; window transitions bypass the feed |
| GWC-27 | Low | P2 | S | GWC-26 | Not started | `AttachInternalEventSubscriber` bypasses the readiness gate; premature attach poisons the distributor |
| GWC-26 | Low | P2 | M | GWC-27 | Done | Alarm monitor attaches its subscriber after SubscribeAlarms; window transitions bypass the feed |
| GWC-27 | Low | P2 | S | GWC-26 | Done | `AttachInternalEventSubscriber` bypasses the readiness gate; premature attach poisons the distributor |
| GWC-28 | Low | P2 | S | GWC-10 (coord, old tracker) | Not started | Gateway→worker envelope `sequence` stamped at creation, not at write |
| GWC-29 | Low | — | S | — | Not started | `Invoke` deep-clones the entire request only to discard the cloned command |
| GWC-30 | Info | — | S | — | Not started | Frame reader allocates a fresh 4-byte length-prefix array per frame |
@@ -164,3 +164,5 @@ Sequence these together rather than piecemeal — several are one change set spa
| 2026-08-07 | **TST-29 → `Done`:** migrated the Phase-5 (orphan-worker reattach) deferred-not-planned governance record and the settled Phase-4 Viewer-default decision from `oldtasks.md` into a new "Session-Resilience Epic Scope" entry in `docs/DesignDecisions.md`; repointed CLAUDE.md and `stillpending.md:7,165` from `oldtasks.md` to `docs/DesignDecisions.md` / `docs/plans/2026-06-15-session-resilience.md.tasks.json`; `git rm oldtasks.md`. The five untracked root docs-review artifacts (`MxAccessGateway-docs-{issues,fixed,final}.md`, `MxGatewayClient-docs-{issues,fixed}.md`) were absent from this worktree — delete from the main working tree separately. |
| 2026-08-07 | **GWC-24 → `Done`** (branch `fix/gwc-24-staging-bound`). `WorkerClient._eventStaging` is now `Channel.CreateBounded` at `2 × EventChannelCapacity` (`Wait`, single reader/writer, no sync continuations); a rejected staging `TryWrite` faults the client `ProtocolViolation` with `QueueOverflow("worker-event-staging")` unless `IsTerminalState()` (shutdown stays a silent drop), so a consumer draining slower than its worker produces dies at a fixed ceiling instead of growing gateway memory. Queue-depth accounting moved from `EnqueueWorkerEventAsync` to `StageWorkerEvent`, so the single gauge reports staged + queued; the timed-write fault (`EventChannelFullModeTimeout` / `QueueOverflow("worker-events")`) is unchanged and still catches the full-stall case first. No new config key — total gateway-side buffering is `3 × MxGateway:Events:QueueCapacity`, derived; coordination with still-open old **GWC-21** (`EventChannelFullModeTimeout` configurability) remains open and was not blocked on. Docs same commit: `GatewayProcessDesign.md` (two overflow faults), `MxAccessWorkerInstanceDesign.md`, `GatewayConfiguration.md`, `Metrics.md`. Tests: `WorkerClientTests.StagingChannelOverflowFaultsWorkerWithoutWaitingForFullModeTimeout` and `.WorkerEventQueueDepthGaugeCountsStagedEvents`; `WorkerClientTests` 22/22 green, `NonWindows.slnx` builds with 0 warnings. |
| 2026-08-07 | **ReplayGap end-to-end cluster (GWC-25 + CLI-35 + CLI-36) → `Done`** on `fix/gwc-25-replaygap-trio`. GWC-25: `SessionEventDistributor.RegisterWithReplay`'s empty-ring branch now reports `oldestAvailableSequence = _highestSequenceSeen + 1` when `gap == true` (still `0` when no gap), so the universal `oldest - 1` resume formula no longer wraps to `ulong.MaxValue` and dead-stream the subscriber; `docs/Sessions.md` documents the empty-ring value. CLI-35: the Python CLI renders a `ReplayGap` as a `{"replayGap": {...}}` row via a new `_event_row` helper instead of crashing in `MessageToDict`. CLI-36: the Go CLI branches on `result.IsReplayGap()` and prints the typed `REPLAY_GAP requested_after=<n> oldest_available=<n>` line / `replayGap` JSON row instead of formatting the library's cleared `Event`. `docs/CrossLanguageSmokeMatrix.md` gained a per-CLI gap-rendering table (one edit covering both client findings). Four new tests as designed (3 × `SessionEventDistributorTests`, `GatewayEndToEndReconnectReplayTests.ReconnectAfterFullAgeEvictionResumesWithSentinelFormula`) plus `test_stream_events_renders_replay_gap` (Python) and `TestRunStreamEventsPrintsReplayGap` (Go); all written red first and each reproducing its defect verbatim. **Deferred:** GWC-25's `ReplayGap.oldest_available_sequence` proto-comment amendment is **not** in this change — it is comment-only but triggers the full five-client regen fan-out, so it lands with the later codegen wave (alongside IPC-23's proto-comment edits) rather than forcing a regen for one sentence. Note for that wave: the fake-worker gateway e2e suite cannot run on the macOS worktree without `TMPDIR` shortened (macOS caps the Unix-domain-socket path backing .NET named pipes at 104 chars; `TMPDIR=/tmp dotnet test …` works and was used here). |
| 2026-08-07 | **GWC-27 → `Done`, GWC-26 → `Done`** (branch `fix/gwc-26-27-alarm-attach`). GWC-27: `GatewaySession.AttachInternalEventSubscriber` now mirrors `AttachEventSubscriber`'s readiness gate under `_syncRoot`, before `EnsureDistributorCreated`, so a premature attach can no longer latch a poisoned distributor. GWC-26: the alarm monitor takes its internal lease directly from the session **before** `SubscribeAlarms` and drains it after the first reconcile; `ISessionManager.ReadAlarmEventsAsync` removed (zero remaining callers); `ApplyReconcile` now broadcasts an `Acknowledge` feed transition for a both-present alarm whose state advanced to `ActiveAcked` (feed-level repair on `AlarmFeedMessage`, not `MxEvent` synthesis). New tests `GatewaySessionTests.AttachInternalEventSubscriberBeforeReadyThrowsAndDoesNotPoisonDistributor` and `GatewayAlarmMonitorAttachOrderTests` (`TransitionsDuringSubscribeWindow_StillReachTheAlarmFeed`, `ApplyReconcileBroadcastsAcknowledgeDelta`); the alarm-monitor fakes now hand the monitor a real Ready `GatewaySession` with a dashboard mirror so the window is actually reproducible. Verification: NonWindows build 0 warnings/0 errors; `GatewayAlarmMonitor` 16 passed, `SessionManagerTests` 38 passed, `GatewaySessionTests` 19 passed, `AlarmFailoverEndToEndTests` 2 passed. |
| 2026-08-07 | Code review of `fix/gwc-26-27-alarm-attach` surfaced a **known pre-existing characteristic, now documented**: the alarm monitor's reconcile-derived feed repairs are **at-least-once, not exactly-once**. A reconcile reads the worker's current state while the matching live transition may still be buffered in the monitor's internal lease, so both broadcast and the duplicates are indistinguishable on the alarm feed (`StreamAlarms` + dashboard alarm hub). This pre-dates GWC-26 — the Raise/Clear presence repair has always had it, since nothing serializes a reconcile pass against the in-flight live stream — so closing it (reconcile/live serialization or transition-timestamp dedup) was ruled out of scope for a P2 fix. Documented instead in `GatewayAlarmMonitor.ApplyReconcile`, `gateway.md`, and `docs/Sessions.md`, with the consumer-side contract stated explicitly (apply transitions idempotently — "set this alarm to this state", never increment/toggle). **Candidate finding for the next review cycle.** |
@@ -10,8 +10,8 @@ This document turns the 2026-07-12 re-review's **new** Gateway Server Core findi
|----|-----|------|-----|-----|--------|-------|
| GWC-24 | Medium | P1 | M | GWC-21 (coord) | Done | Unbounded event staging channel: sustained slow drain grows memory silently and invisibly |
| GWC-25 | Medium | P0 | S | CLI-35/36 (coord) | Done | Empty-ring ReplayGap sentinel carries `oldest_available_sequence = 0`, dead-streaming a compliant client |
| GWC-26 | Low | P2 | M | GWC-27 | Not started | Alarm monitor attaches its subscriber after SubscribeAlarms; window transitions bypass the feed, missed Acknowledge never repaired |
| GWC-27 | Low | P2 | S | GWC-26 | Not started | `AttachInternalEventSubscriber` bypasses the readiness gate; premature attach poisons the distributor permanently |
| GWC-26 | Low | P2 | M | GWC-27 | Done | Alarm monitor attaches its subscriber after SubscribeAlarms; window transitions bypass the feed, missed Acknowledge never repaired |
| GWC-27 | Low | P2 | S | GWC-26 | Done | `AttachInternalEventSubscriber` bypasses the readiness gate; premature attach poisons the distributor permanently |
| GWC-28 | Low | P2 | S | GWC-10 (coord) | Not started | Gateway→worker envelope `sequence` stamped at creation, not at write — non-monotonic on the wire under concurrent invokes |
| GWC-29 | Low | — | S | — | Not started | `Invoke` deep-clones the entire request only to discard the cloned command |
| GWC-30 | Info | — | S | — | Not started | Frame reader allocates a fresh 4-byte length-prefix array per frame |
+7 -1
View File
@@ -197,7 +197,13 @@ Event streaming uses `AttachEventSubscriber` which returns a disposable lease. W
`FailFast` event backpressure faults the whole session only in single-subscriber mode; in multi-subscriber mode it degrades to a per-subscriber disconnect so one slow consumer never faults a session shared by others. The session passes its mode to the `SessionEventDistributor` at construction, so this decision is made on the fixed mode rather than a live subscriber-count snapshot.
The single worker event channel has exactly one direct reader: the `SessionEventDistributor` pump (`MapWorkerEventsAsync`). Both gateway-owned internal consumers — the dashboard mirror and the central alarm monitor — attach as distributor subscribers rather than draining the worker channel themselves. `GatewaySession.AttachInternalEventSubscriber` mirrors the dashboard-mirror lease (`isInternal: true`): the alarm monitor's `SessionManager.ReadAlarmEventsAsync` registers one so it consumes the same mapped `MxEvent`s the pump fans to every subscriber, without counting against `MaxEventSubscribersPerSession` and without a slow reconcile faulting the session. This is what keeps the alarm feed and the dashboard from splitting the stream between two raw drains (which would silently lose Acknowledge and provider-mode transitions); the worker channel is single-reader and a second `WorkerClient.ReadEventsAsync` consumer throws so a regression fails loudly.
The single worker event channel has exactly one direct reader: the `SessionEventDistributor` pump (`MapWorkerEventsAsync`). Both gateway-owned internal consumers — the dashboard mirror and the central alarm monitor — attach as distributor subscribers rather than draining the worker channel themselves. `GatewaySession.AttachInternalEventSubscriber` mirrors the dashboard-mirror lease (`isInternal: true`): the alarm monitor calls it directly on its session so it consumes the same mapped `MxEvent`s the pump fans to every subscriber, without counting against `MaxEventSubscribersPerSession` and without a slow reconcile faulting the session. This is what keeps the alarm feed and the dashboard from splitting the stream between two raw drains (which would silently lose Acknowledge and provider-mode transitions); the worker channel is single-reader and a second `WorkerClient.ReadEventsAsync` consumer throws so a regression fails loudly.
The monitor takes that lease **before** it issues `SubscribeAlarms`, so no transition window is missed: the pump has been running since `MarkReady` started the dashboard mirror, and the distributor only fans to subscribers registered at fan-out time, so a lease taken after the subscribe + first-reconcile round trips would lose every transition raised inside that window. Transitions arriving while the monitor subscribes and reconciles buffer in the lease's bounded channel and are applied after the reconcile, which is order-safe because a live transition can update an alarm the reconciled snapshot already holds.
The repair transitions the monitor's reconcile broadcasts on the alarm feed (Raise/Clear presence deltas and the acked-state delta) are **at-least-once**: a reconcile reads the worker's current state while the matching live transition may still be buffered in the lease, so both can be broadcast and the duplicates are indistinguishable — alarm-feed consumers must apply transitions idempotently. This is a property of the reconcile design, not of the buffering above.
`AttachInternalEventSubscriber` enforces the same readiness gate as `AttachEventSubscriber` — a session (or worker) that is not `Ready` throws `SessionNotReady` *before* the distributor is constructed. A premature attach would otherwise start the pump against a source that throws, completing every subscriber with that error and latching the distributor for the session's whole lifetime.
Sessions open with `MxGateway:Sessions:DefaultLeaseSeconds` (default 1800) added to the open timestamp. Unary client activity refreshes the lease by the same duration. `ExtendLease` and `IsLeaseExpired` cooperate with `SessionManager.CloseExpiredLeasesAsync`, which iterates a registry snapshot and closes any session whose lease has expired with `LeaseExpiredReason`. `SessionLeaseMonitorHostedService` runs that sweep every `MxGateway:Sessions:LeaseSweepIntervalSeconds` seconds (default 30).
+23
View File
@@ -154,6 +154,29 @@ session. The worker event channel is single-reader and asserts it (a second
`WorkerClient.ReadEventsAsync` consumer throws), so a regression cannot silently
split the event stream between two drains.
The monitor takes that internal lease **before** it sends `SubscribeAlarms`, and
drains it after the first reconcile. The pump is already running by then (the
dashboard mirror starts it at `MarkReady`) and the distributor fans only to
subscribers registered at fan-out time, so attaching after the subscribe +
reconcile round trips would drop every transition raised in that window —
including an `Acknowledge`, which the presence-only reconcile deltas would never
repair. Transitions arriving during the window buffer in the lease's bounded
channel instead. As a backstop for any window this ordering cannot cover (worker
restart, internal-subscriber overflow disconnect), a reconcile that finds a known
alarm now reported `ActiveAcked` broadcasts an `Acknowledge` transition on the
alarm feed. That is a feed-level repair rebuilt from the worker's own snapshot on
the `StreamAlarms` surface — it is not an `MxEvent` and never reaches
`StreamEvents`, so the "never synthesize events" rule is untouched.
**Feed repair transitions are at-least-once, not exactly-once.** A reconcile reads
the worker's current state while the corresponding live transition may still be
buffered in the monitor's lease, so both can broadcast and the two are
indistinguishable on the feed. This applies to the acked-state delta and equally
to the older Raise/Clear presence repair: nothing serializes a reconcile pass
against the in-flight live stream. Alarm-feed consumers (`StreamAlarms` clients
and the dashboard alarm hub) must apply transitions idempotently — treat one as
"set this alarm to this state", never as an increment or a toggle.
### Alarm providers and failover
The alarm feed has two providers, both implemented worker-side:
@@ -210,6 +210,17 @@ public sealed class GatewayAlarmMonitor : BackgroundService, IGatewayAlarmServic
try
{
// Attach the internal distributor subscriber BEFORE subscribing (GWC-26). The pump
// has been running since MarkReady started the dashboard mirror, and the distributor
// only fans to subscribers registered at the time of the fan-out, so a subscriber
// taken after SubscribeAlarms + the first reconcile would silently lose every
// transition raised inside that two-round-trip window — and a missed Acknowledge is
// never repaired by the presence-only reconcile deltas. Transitions arriving while we
// subscribe and reconcile simply buffer in this lease's bounded channel; if it ever
// overflowed, the internal subscriber is disconnected (it never faults the session),
// the enumeration below ends, and the supervisor loop restarts the lifecycle.
using IEventSubscriberLease alarmLease = session.AttachInternalEventSubscriber();
await SubscribeAlarmsAsync(session.SessionId, subscription, stoppingToken).ConfigureAwait(false);
await ReconcileAsync(session.SessionId, stoppingToken).ConfigureAwait(false);
@@ -228,9 +239,11 @@ public sealed class GatewayAlarmMonitor : BackgroundService, IGatewayAlarmServic
// Consume mapped MxEvents through the session's single distributor pump (as an
// internal, non-counted subscriber) rather than opening a second raw drain of the
// worker event channel — a second drain would split events with the dashboard
// mirror pump and silently lose Acknowledge/mode-change transitions.
await foreach (MxEvent mxEvent in _sessionManager
.ReadAlarmEventsAsync(session.SessionId, linked.Token)
// mirror pump and silently lose Acknowledge/mode-change transitions. The lease was
// taken above, before SubscribeAlarms; draining it only now is order-safe because
// ApplyTransition handles alarms the reconcile already placed in the cache.
await foreach (MxEvent mxEvent in alarmLease.Reader
.ReadAllAsync(linked.Token)
.ConfigureAwait(false))
{
if (mxEvent is { BodyCase: MxEvent.BodyOneofCase.OnAlarmTransition }
@@ -511,6 +524,20 @@ public sealed class GatewayAlarmMonitor : BackgroundService, IGatewayAlarmServic
// Replaces the cache with the worker's authoritative snapshot, broadcasting
// a synthetic transition for any alarm the live stream missed.
//
// Repair scope (GWC-26): presence deltas (Clear/Raise) plus acked-state deltas. These are
// ALARM FEED transitions (AlarmFeedMessage on the StreamAlarms/dashboard surface), rebuilt
// from the worker's own authoritative snapshot to repair what the live feed missed. They are
// not MxEvents and never reach the gRPC StreamEvents path, so this feed-level repair does not
// breach the "never synthesize events" rule, which governs MxEvent emission.
//
// Delivery semantics: feed repair transitions are AT-LEAST-ONCE, not exactly-once. A reconcile
// reads the worker's current state while the corresponding live transition may still be
// buffered in the alarm lease's channel; both then broadcast, and the two are indistinguishable
// on the feed. This is inherent to the reconcile design and pre-dates the acked-state delta
// (the Raise/Clear repair has always had it), since nothing serializes a reconcile against the
// in-flight live stream. Consumers must therefore treat alarm state idempotently — apply a
// transition as "set the alarm to this state", never as an increment or a toggle.
private void ApplyReconcile(IEnumerable<ActiveAlarmSnapshot> snapshots)
{
Dictionary<string, ActiveAlarmSnapshot> next = new(StringComparer.Ordinal);
@@ -536,12 +563,23 @@ public sealed class GatewayAlarmMonitor : BackgroundService, IGatewayAlarmServic
foreach (KeyValuePair<string, ActiveAlarmSnapshot> incoming in next)
{
if (!_alarms.ContainsKey(incoming.Key))
if (!_alarms.TryGetValue(incoming.Key, out ActiveAlarmSnapshot? existing))
{
Broadcast(
new AlarmFeedMessage { Transition = TransitionFromSnapshot(incoming.Value, AlarmTransitionKind.Raise) },
incoming.Key);
}
else if (existing.CurrentState != incoming.Value.CurrentState
&& incoming.Value.CurrentState == AlarmConditionState.ActiveAcked)
{
// The alarm was already known but the worker now reports it acknowledged: the
// live Acknowledge transition never reached the feed. Without this the acked
// state is absorbed silently by the snapshot replace below and subscribers show
// the alarm unacked until it clears.
Broadcast(
new AlarmFeedMessage { Transition = TransitionFromSnapshot(incoming.Value, AlarmTransitionKind.Acknowledge) },
incoming.Key);
}
}
_alarms.Clear();
@@ -548,10 +548,37 @@ public sealed class GatewaySession
/// <c>MaxEventSubscribersPerSession</c> accounting and out of the single-subscriber
/// overflow-fault path, so a slow alarm reconcile can never fault the session — it only
/// disconnects this internal subscriber.
/// <para>
/// Gated on readiness exactly like <see cref="AttachEventSubscriber"/>: attaching
/// before the session and its worker are <c>Ready</c> throws
/// <see cref="SessionManagerException"/> with
/// <see cref="SessionManagerErrorCode.SessionNotReady"/>.
/// </para>
/// </remarks>
/// <returns>The internal subscriber's lease; dispose it to unregister.</returns>
/// <exception cref="SessionManagerException">
/// The session or its worker client is not <c>Ready</c>.
/// </exception>
public IEventSubscriberLease AttachInternalEventSubscriber()
{
// Readiness gate, mirroring AttachEventSubscriber (GWC-27). It must run BEFORE
// EnsureDistributorCreated: a premature attach would construct the distributor and start
// its pump against a not-yet-Ready worker, the pump source would throw SessionNotReady,
// PumpAsync would complete every subscriber with that error and latch _completed, and
// _eventDistributorStarted is never reset — so the session would reach Ready with
// permanently dead event streaming, silently, for the rest of its life. Failing loudly
// here keeps that state unreachable. The check is under _syncRoot and the distributor
// calls stay outside it, matching AttachEventSubscriber's lock discipline.
lock (_syncRoot)
{
if (_state != SessionState.Ready || _workerClient?.State != WorkerClientState.Ready)
{
throw new SessionManagerException(
SessionManagerErrorCode.SessionNotReady,
$"Session {SessionId} is not ready for event streaming. Current state is {_state}.");
}
}
// Same sequence StartDashboardMirror uses: create the distributor (claiming the pump
// start if we are first), register the internal subscriber BEFORE the pump starts so a
// subscriber is always present at pump start, then start the pump if requested.
@@ -43,18 +43,6 @@ public interface ISessionManager
string sessionId,
CancellationToken cancellationToken);
/// <summary>
/// Reads mapped events for the central alarm monitor by attaching an internal
/// (non-counted) distributor subscriber, so the alarm feed shares the one worker-event
/// pump instead of opening a second raw drain of the single worker event channel.
/// </summary>
/// <param name="sessionId">Identifier of the session.</param>
/// <param name="cancellationToken">Token to cancel the asynchronous operation.</param>
/// <returns>The mapped <see cref="MxEvent"/>s fanned by the session's distributor.</returns>
IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken);
/// <summary>Closes a session and terminates its worker process.</summary>
/// <param name="sessionId">Identifier of the session to close.</param>
/// <param name="cancellationToken">Token to cancel the asynchronous operation.</param>
@@ -1,5 +1,4 @@
using System.Diagnostics.CodeAnalysis;
using System.Runtime.CompilerServices;
using System.Security.Cryptography;
using Google.Protobuf.WellKnownTypes;
using Microsoft.Extensions.Logging;
@@ -191,22 +190,6 @@ public sealed class SessionManager : ISessionManager
return session.ReadEventsAsync(cancellationToken);
}
/// <inheritdoc />
public async IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
[EnumeratorCancellation] CancellationToken cancellationToken)
{
GatewaySession session = GetRequiredSession(sessionId);
using IEventSubscriberLease lease = session.AttachInternalEventSubscriber();
await foreach (MxEvent mxEvent in lease.Reader
.ReadAllAsync(cancellationToken)
.ConfigureAwait(false))
{
yield return mxEvent;
}
}
/// <inheritdoc />
public async Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -1,5 +1,4 @@
using System.Diagnostics.CodeAnalysis;
using System.Runtime.CompilerServices;
using System.Threading.Channels;
using Google.Protobuf.WellKnownTypes;
using Microsoft.Extensions.Logging.Abstractions;
@@ -8,6 +7,7 @@ using ZB.MOM.WW.MxGateway.Server.Alarms;
using ZB.MOM.WW.MxGateway.Server.Configuration;
using ZB.MOM.WW.MxGateway.Server.Metrics;
using ZB.MOM.WW.MxGateway.Server.Sessions;
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
namespace ZB.MOM.WW.MxGateway.Tests.Alarms;
@@ -420,6 +420,9 @@ public sealed class AlarmFailoverEndToEndTests
string? ownerKeyId,
CancellationToken cancellationToken)
{
// The monitor attaches its internal subscriber directly on this session, so the
// session has to be a genuinely Ready one with a worker client feeding the
// distributor pump — EmitEvent writes into that worker's event stream.
GatewaySession session = new(
Guid.NewGuid().ToString("N"),
"Galaxy",
@@ -432,6 +435,8 @@ public sealed class AlarmFailoverEndToEndTests
TimeSpan.FromSeconds(30),
TimeSpan.FromSeconds(30),
DateTimeOffset.UtcNow);
session.AttachWorkerClient(new ChannelWorkerClient(session.SessionId, _events.Reader));
session.MarkReady();
return Task.FromResult(session);
}
@@ -460,29 +465,9 @@ public sealed class AlarmFailoverEndToEndTests
}
/// <inheritdoc />
public async IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
public IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
string sessionId,
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
{
yield return workerEvent;
}
}
/// <inheritdoc />
public async IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
{
if (workerEvent.Event is not null)
{
yield return workerEvent.Event;
}
}
}
CancellationToken cancellationToken) => throw new NotSupportedException();
/// <inheritdoc />
public bool TryGetSession(string sessionId, [MaybeNullWhen(false)] out GatewaySession session)
@@ -0,0 +1,468 @@
using System.Diagnostics.CodeAnalysis;
using System.Threading.Channels;
using Google.Protobuf.WellKnownTypes;
using Microsoft.Extensions.Logging.Abstractions;
using ZB.MOM.WW.MxGateway.Contracts.Proto;
using ZB.MOM.WW.MxGateway.Server.Alarms;
using ZB.MOM.WW.MxGateway.Server.Configuration;
using ZB.MOM.WW.MxGateway.Server.Grpc;
using ZB.MOM.WW.MxGateway.Server.Metrics;
using ZB.MOM.WW.MxGateway.Server.Sessions;
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
namespace ZB.MOM.WW.MxGateway.Tests.Alarms;
/// <summary>
/// GWC-26 regression tests for the alarm monitor's startup ordering. Unlike the
/// sibling alarm-monitor tests, the session manager here hands the monitor a REAL
/// <see cref="GatewaySession"/> that is driven to Ready with a dashboard mirror, so the
/// distributor pump is already running when the monitor attaches — the production
/// condition under which a late internal subscriber silently misses everything the pump
/// already fanned.
/// </summary>
public sealed class GatewayAlarmMonitorAttachOrderTests
{
private const string AlarmReference = "Galaxy!Area.Tank01.Hi";
private static readonly TimeSpan WaitTimeout = TimeSpan.FromSeconds(30);
/// <summary>
/// The monitor must take its internal distributor lease BEFORE issuing
/// <c>SubscribeAlarms</c>, so transitions the worker emits during the
/// subscribe + first-reconcile window buffer in the lease instead of being fanned to a
/// subscriber set the monitor has not joined yet. The test parks the monitor inside its
/// <c>SubscribeAlarms</c> round trip, emits a Raise and an Acknowledge, waits until the
/// dashboard mirror proves the pump has already fanned both, and only then releases the
/// monitor: with a late attach both transitions are lost to the alarm feed forever.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task TransitionsDuringSubscribeWindow_StillReachTheAlarmFeed()
{
using GatewayMetrics metrics = new();
await using FakeSessionManager sessions = new();
sessions.HoldSubscribeUntilReleased();
using GatewayAlarmMonitor monitor = CreateMonitor(sessions, metrics);
using CancellationTokenSource cts = new();
await monitor.StartAsync(cts.Token);
await sessions.WaitForSubscribeStartAsync(WaitTimeout);
// A live feed subscriber, drained past its baseline ProviderStatus so it is registered
// before any window transition is broadcast.
List<AlarmFeedMessage> received = [];
TaskCompletionSource baselineReceived = new(TaskCreationOptions.RunContinuationsAsynchronously);
using CancellationTokenSource streamCts = new();
Task reader = ReadFeedAsync(monitor, received, baselineReceived, streamCts.Token);
await baselineReceived.Task.WaitAsync(WaitTimeout);
// The window: the worker reports a raise and an acknowledge while the monitor is still
// waiting for its SubscribeAlarms reply.
sessions.EmitEvent(Transition(1, AlarmTransitionKind.Raise));
sessions.EmitEvent(Transition(2, AlarmTransitionKind.Acknowledge));
// The dashboard mirror is an independent distributor subscriber: once it has both
// events, the pump has provably already fanned them, so a subscriber that registers
// after this point can never receive them.
await WaitUntilAsync(() => sessions.Broadcaster.Captures.Count == 2, WaitTimeout);
sessions.ReleaseSubscribe();
AlarmFeedMessage raise = await WaitForAsync(
received,
m => m.PayloadCase == AlarmFeedMessage.PayloadOneofCase.Transition
&& m.Transition.TransitionKind == AlarmTransitionKind.Raise,
WaitTimeout);
AlarmFeedMessage acknowledge = await WaitForAsync(
received,
m => m.PayloadCase == AlarmFeedMessage.PayloadOneofCase.Transition
&& m.Transition.TransitionKind == AlarmTransitionKind.Acknowledge,
WaitTimeout);
Assert.Equal(AlarmReference, raise.Transition.AlarmFullReference);
Assert.Equal(AlarmReference, acknowledge.Transition.AlarmFullReference);
await streamCts.CancelAsync();
await reader;
await cts.CancelAsync();
await monitor.StopAsync(CancellationToken.None);
}
/// <summary>
/// Defense in depth for any window the attach reorder cannot cover (worker restart,
/// internal-subscriber overflow disconnect): when a reconcile snapshot reports an alarm
/// the cache already holds but with the state advanced to
/// <see cref="AlarmConditionState.ActiveAcked"/>, the monitor broadcasts an
/// <see cref="AlarmTransitionKind.Acknowledge"/> feed transition. Before the fix the
/// acked state was absorbed silently by the snapshot replace, leaving live subscribers
/// showing the alarm unacked until it cleared.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task ApplyReconcileBroadcastsAcknowledgeDelta()
{
using GatewayMetrics metrics = new();
await using FakeSessionManager sessions = new();
using GatewayAlarmMonitor monitor = CreateMonitor(sessions, metrics);
using CancellationTokenSource cts = new();
await monitor.StartAsync(cts.Token);
await sessions.WaitForSubscribeStartAsync(WaitTimeout);
// Seed the cache with an active (unacked) alarm through a reconcile rather than a live
// transition: a buffered live transition could still be in flight when the cache first
// shows the alarm, and would then broadcast to the reader below and pollute the
// exactly-one-transition assertion. A provider-mode event forces the reconcile
// immediately, so the test never waits on the periodic timer.
sessions.SetReconcileSnapshot(Snapshot(AlarmConditionState.Active));
sessions.EmitEvent(ProviderModeProbe(1));
await WaitUntilAsync(
() => monitor.CurrentAlarms.Any(alarm => alarm.AlarmFullReference == AlarmReference
&& alarm.CurrentState == AlarmConditionState.Active),
WaitTimeout);
// Subscribe AFTER the seed: this reader's snapshot carries the alarm, so every
// transition it observes from here on is a reconcile-derived broadcast.
List<AlarmFeedMessage> received = [];
TaskCompletionSource snapshotComplete = new(TaskCreationOptions.RunContinuationsAsynchronously);
using CancellationTokenSource streamCts = new();
Task reader = ReadFeedAsync(monitor, received, snapshotComplete, streamCts.Token, untilSnapshotComplete: true);
await snapshotComplete.Task.WaitAsync(WaitTimeout);
// The worker now reports the same alarm acknowledged. A provider-mode event forces an
// immediate reconcile pass so the test does not wait on the periodic timer.
sessions.SetReconcileSnapshot(Snapshot(AlarmConditionState.ActiveAcked));
sessions.EmitEvent(ProviderModeProbe(2));
AlarmFeedMessage acknowledge = await WaitForAsync(
received,
m => m.PayloadCase == AlarmFeedMessage.PayloadOneofCase.Transition
&& m.Transition.TransitionKind == AlarmTransitionKind.Acknowledge,
WaitTimeout);
Assert.Equal(AlarmReference, acknowledge.Transition.AlarmFullReference);
lock (received)
{
AlarmFeedMessage[] transitions = received
.Where(m => m.PayloadCase == AlarmFeedMessage.PayloadOneofCase.Transition)
.ToArray();
AlarmFeedMessage single = Assert.Single(transitions);
Assert.Equal(AlarmTransitionKind.Acknowledge, single.Transition.TransitionKind);
}
await streamCts.CancelAsync();
await reader;
await cts.CancelAsync();
await monitor.StopAsync(CancellationToken.None);
}
private static GatewayAlarmMonitor CreateMonitor(FakeSessionManager sessions, GatewayMetrics metrics)
{
AlarmsOptions options = new()
{
Enabled = true,
SubscriptionExpression = @"\\NODE\Galaxy!Area",
};
return new GatewayAlarmMonitor(
sessions,
new StubWatchListResolver(),
metrics,
Microsoft.Extensions.Options.Options.Create(new GatewayOptions { Alarms = options }),
NullLogger<GatewayAlarmMonitor>.Instance);
}
// Drains the monitor's feed into received, signalling gate on the first message (the
// baseline ProviderStatus) or, when untilSnapshotComplete is set, on SnapshotComplete.
private static Task ReadFeedAsync(
GatewayAlarmMonitor monitor,
List<AlarmFeedMessage> received,
TaskCompletionSource gate,
CancellationToken cancellationToken,
bool untilSnapshotComplete = false)
{
return Task.Run(
async () =>
{
try
{
await foreach (AlarmFeedMessage message in monitor.StreamAsync(null, cancellationToken))
{
bool opensGate = !untilSnapshotComplete
|| message.PayloadCase == AlarmFeedMessage.PayloadOneofCase.SnapshotComplete;
// Record only what arrives AFTER the gate opened: everything up to and
// including the gate message is this subscriber's snapshot preamble, not
// a live broadcast, so assertions stay about broadcasts alone.
lock (received)
{
if (gate.Task.IsCompleted)
{
received.Add(message);
}
}
if (opensGate)
{
gate.TrySetResult();
}
}
}
catch (OperationCanceledException)
{
// Expected when the test cancels the stream.
}
},
CancellationToken.None);
}
private static MxEvent Transition(ulong sequence, AlarmTransitionKind kind) => new()
{
Family = MxEventFamily.OnAlarmTransition,
WorkerSequence = sequence,
OnAlarmTransition = new OnAlarmTransitionEvent
{
AlarmFullReference = AlarmReference,
SourceObjectReference = "Tank01",
AlarmTypeName = "AnalogLimitAlarm.Hi",
TransitionKind = kind,
Severity = 500,
SourceProvider = AlarmProviderMode.Alarmmgr,
TransitionTimestamp = Timestamp.FromDateTimeOffset(DateTimeOffset.UtcNow),
},
};
// A no-op provider-mode event. The monitor forces an immediate reconcile after every one,
// which is how these tests drive a reconcile pass without waiting on the periodic timer.
private static MxEvent ProviderModeProbe(ulong sequence) => new()
{
Family = MxEventFamily.OnAlarmProviderModeChanged,
WorkerSequence = sequence,
OnAlarmProviderModeChanged = new OnAlarmProviderModeChangedEvent
{
Mode = AlarmProviderMode.Alarmmgr,
Reason = "probe",
At = Timestamp.FromDateTimeOffset(DateTimeOffset.UtcNow),
},
};
private static ActiveAlarmSnapshot Snapshot(AlarmConditionState state) => new()
{
AlarmFullReference = AlarmReference,
SourceObjectReference = "Tank01",
AlarmTypeName = "AnalogLimitAlarm.Hi",
CurrentState = state,
Severity = 500,
SourceProvider = AlarmProviderMode.Alarmmgr,
};
private static async Task<AlarmFeedMessage> WaitForAsync(
List<AlarmFeedMessage> received,
Func<AlarmFeedMessage, bool> predicate,
TimeSpan timeout)
{
DateTime deadline = DateTime.UtcNow + timeout;
while (DateTime.UtcNow < deadline)
{
lock (received)
{
AlarmFeedMessage? match = received.FirstOrDefault(predicate);
if (match is not null)
{
return match;
}
}
await Task.Delay(25);
}
throw new TimeoutException("No matching AlarmFeedMessage was received in time.");
}
private static async Task WaitUntilAsync(Func<bool> condition, TimeSpan timeout)
{
DateTime deadline = DateTime.UtcNow + timeout;
while (DateTime.UtcNow < deadline)
{
if (condition())
{
return;
}
await Task.Delay(25);
}
throw new TimeoutException("Condition was not met in time.");
}
/// <summary><see cref="IAlarmWatchListResolver"/> that resolves an empty watch-list.</summary>
private sealed class StubWatchListResolver : IAlarmWatchListResolver
{
/// <inheritdoc />
public Task<IReadOnlyList<AlarmSubtagTarget>> ResolveAsync(
AlarmsOptions options,
CancellationToken cancellationToken = default) =>
Task.FromResult<IReadOnlyList<AlarmSubtagTarget>>([]);
}
/// <summary>
/// Session manager that hands the monitor a real <see cref="GatewaySession"/> driven to
/// Ready with a dashboard mirror, so the distributor pump is running before the monitor
/// attaches. <see cref="EmitEvent"/> pushes worker events through the fake worker client
/// into that pump, exactly as a live worker would.
/// </summary>
private sealed class FakeSessionManager : ISessionManager, IAsyncDisposable
{
private readonly Channel<WorkerEvent> _events = Channel.CreateUnbounded<WorkerEvent>();
private readonly TaskCompletionSource _subscribeStarted =
new(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly object _sync = new();
private TaskCompletionSource _subscribeGate = CreateReleasedGate();
private ActiveAlarmSnapshot[] _reconcileSnapshot = [];
private GatewaySession? _session;
/// <summary>Dashboard mirror attached to the session; proves what the pump has fanned.</summary>
public RecordingDashboardEventBroadcaster Broadcaster { get; } = new();
/// <summary>Re-arms the gate so <c>SubscribeAlarms</c> parks until <see cref="ReleaseSubscribe"/>.</summary>
public void HoldSubscribeUntilReleased() =>
_subscribeGate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
/// <summary>Releases a gate armed by <see cref="HoldSubscribeUntilReleased"/>.</summary>
public void ReleaseSubscribe() => _subscribeGate.TrySetResult();
/// <summary>Completes once the monitor's <c>SubscribeAlarms</c> command has arrived.</summary>
/// <param name="timeout">The maximum time to wait.</param>
/// <returns>A task that completes when the command arrives.</returns>
public Task WaitForSubscribeStartAsync(TimeSpan timeout) => _subscribeStarted.Task.WaitAsync(timeout);
/// <summary>Sets the active-alarm snapshot every <c>QueryActiveAlarms</c> reconcile returns.</summary>
/// <param name="snapshots">The snapshots to report.</param>
public void SetReconcileSnapshot(params ActiveAlarmSnapshot[] snapshots)
{
lock (_sync)
{
_reconcileSnapshot = snapshots;
}
}
/// <summary>Pushes a worker event into the session's distributor pump.</summary>
/// <param name="mxEvent">The event to push.</param>
public void EmitEvent(MxEvent mxEvent) =>
_events.Writer.TryWrite(new WorkerEvent { Event = mxEvent });
/// <inheritdoc />
public Task<GatewaySession> OpenSessionAsync(
SessionOpenRequest request,
string? clientIdentity,
string? ownerKeyId,
CancellationToken cancellationToken)
{
GatewaySession session = new(
sessionId: "session-alarm-attach-order",
backendName: "Galaxy",
pipeName: "mxaccess-gateway-1-session-alarm-attach-order",
nonce: "nonce",
clientIdentity: clientIdentity,
ownerKeyId: ownerKeyId,
clientSessionName: request.ClientSessionName,
clientCorrelationId: request.ClientCorrelationId,
commandTimeout: TimeSpan.FromSeconds(30),
startupTimeout: TimeSpan.FromSeconds(30),
shutdownTimeout: TimeSpan.FromSeconds(30),
leaseDuration: TimeSpan.FromMinutes(30),
openedAt: DateTimeOffset.UtcNow,
eventStreaming: new SessionEventStreaming(
new MxAccessGrpcMapper(),
new EventOptions { QueueCapacity = 64 },
NullLogger<SessionEventDistributor>.Instance,
TimeProvider.System,
new GatewayMetrics(),
Broadcaster));
session.AttachWorkerClient(new ChannelWorkerClient(session.SessionId, _events.Reader));
// MarkReady starts the dashboard mirror, and with it the distributor pump — the
// production precondition this regression depends on.
session.MarkReady();
_session = session;
return Task.FromResult(session);
}
/// <inheritdoc />
public async Task<WorkerCommandReply> InvokeAsync(
string sessionId,
WorkerCommand command,
CancellationToken cancellationToken)
{
MxCommandReply reply = new()
{
ProtocolStatus = new ProtocolStatus { Code = ProtocolStatusCode.Ok },
};
switch (command.Command?.Kind)
{
case MxCommandKind.SubscribeAlarms:
_subscribeStarted.TrySetResult();
await _subscribeGate.Task.WaitAsync(cancellationToken).ConfigureAwait(false);
break;
case MxCommandKind.QueryActiveAlarms:
QueryActiveAlarmsReplyPayload payload = new();
lock (_sync)
{
payload.Snapshots.AddRange(_reconcileSnapshot.Select(snapshot => snapshot.Clone()));
}
reply.QueryActiveAlarms = payload;
break;
}
return new WorkerCommandReply { Reply = reply };
}
/// <inheritdoc />
public IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
string sessionId,
CancellationToken cancellationToken) => throw new NotSupportedException();
/// <inheritdoc />
public bool TryGetSession(string sessionId, [MaybeNullWhen(false)] out GatewaySession session)
{
session = _session;
return session is not null;
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(string sessionId, CancellationToken cancellationToken)
{
_events.Writer.TryComplete();
return Task.FromResult(new SessionCloseResult(sessionId, SessionState.Closed, AlreadyClosed: false));
}
/// <inheritdoc />
public Task<SessionCloseResult> KillWorkerAsync(string sessionId, string reason, CancellationToken cancellationToken) =>
Task.FromResult(new SessionCloseResult(sessionId, SessionState.Closed, AlreadyClosed: false));
/// <inheritdoc />
public Task<int> CloseExpiredLeasesAsync(DateTimeOffset now, CancellationToken cancellationToken) =>
Task.FromResult(0);
/// <inheritdoc />
public Task ShutdownAsync(CancellationToken cancellationToken) => Task.CompletedTask;
/// <summary>Disposes the session the fake handed out.</summary>
/// <returns>A task that represents the asynchronous operation.</returns>
public async ValueTask DisposeAsync()
{
_events.Writer.TryComplete();
if (_session is not null)
{
await _session.DisposeAsync().ConfigureAwait(false);
}
}
private static TaskCompletionSource CreateReleasedGate()
{
TaskCompletionSource gate = new(TaskCreationOptions.RunContinuationsAsynchronously);
gate.SetResult();
return gate;
}
}
}
@@ -1,6 +1,5 @@
using System.Diagnostics.CodeAnalysis;
using System.Diagnostics.Metrics;
using System.Runtime.CompilerServices;
using System.Threading.Channels;
using Google.Protobuf.WellKnownTypes;
using Microsoft.Extensions.Logging.Abstractions;
@@ -10,6 +9,7 @@ using ZB.MOM.WW.MxGateway.Server.Alarms;
using ZB.MOM.WW.MxGateway.Server.Configuration;
using ZB.MOM.WW.MxGateway.Server.Metrics;
using ZB.MOM.WW.MxGateway.Server.Sessions;
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
namespace ZB.MOM.WW.MxGateway.Tests.Alarms;
@@ -733,6 +733,9 @@ public sealed class GatewayAlarmMonitorProviderModeTests
string? ownerKeyId,
CancellationToken cancellationToken)
{
// The monitor attaches its internal subscriber directly on this session, so the
// session has to be a genuinely Ready one with a worker client feeding the
// distributor pump — EmitEvent writes into that worker's event stream.
GatewaySession session = new(
Guid.NewGuid().ToString("N"),
"Galaxy",
@@ -745,6 +748,8 @@ public sealed class GatewayAlarmMonitorProviderModeTests
TimeSpan.FromSeconds(30),
TimeSpan.FromSeconds(30),
DateTimeOffset.UtcNow);
session.AttachWorkerClient(new ChannelWorkerClient(session.SessionId, _events.Reader));
session.MarkReady();
return Task.FromResult(session);
}
@@ -773,29 +778,9 @@ public sealed class GatewayAlarmMonitorProviderModeTests
}
/// <inheritdoc />
public async IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
public IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
string sessionId,
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
{
yield return workerEvent;
}
}
/// <inheritdoc />
public async IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (WorkerEvent workerEvent in _events.Reader.ReadAllAsync(cancellationToken))
{
if (workerEvent.Event is not null)
{
yield return workerEvent.Event;
}
}
}
CancellationToken cancellationToken) => throw new NotSupportedException();
/// <inheritdoc />
public bool TryGetSession(string sessionId, [MaybeNullWhen(false)] out GatewaySession session)
@@ -378,14 +378,6 @@ public sealed class DashboardSessionAdminServiceTests
throw new NotSupportedException();
}
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken)
{
throw new NotSupportedException();
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -755,14 +755,6 @@ public sealed class EventStreamServiceTests
return _sessions[sessionId].ReadEventsAsync(cancellationToken);
}
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken)
{
throw new NotSupportedException();
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -948,14 +948,6 @@ public sealed class MxAccessGatewayServiceConstraintTests
}
}
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken)
{
throw new NotSupportedException();
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -616,14 +616,6 @@ public sealed class MxAccessGatewayServiceTests
}
}
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken)
{
throw new NotSupportedException();
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -352,11 +352,6 @@ public sealed class GatewaySessionDashboardMirrorTests
string sessionId,
CancellationToken cancellationToken) => session.ReadEventsAsync(cancellationToken);
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken) => throw new NotSupportedException();
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -668,6 +668,81 @@ public sealed class GatewaySessionTests
Assert.Equal(SessionState.Ready, session.State);
}
/// <summary>
/// GWC-27: <see cref="GatewaySession.AttachInternalEventSubscriber"/> must refuse to
/// attach before the session is Ready. Without the gate the attach would construct and
/// start the distributor against a not-yet-Ready worker; the pump source throws
/// <c>SessionNotReady</c>, every subscriber is completed with that error, and the
/// distributor latches — leaving a session that reaches Ready with permanently dead
/// event streaming. The second half of the test is the load-bearing one: after the
/// failed attach the session still streams live events, proving the distributor was
/// never created or started.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task AttachInternalEventSubscriberBeforeReadyThrowsAndDoesNotPoisonDistributor()
{
FakeWorkerClient workerClient = new();
workerClient.Events.Add(new WorkerEvent
{
Event = new MxEvent { Family = MxEventFamily.OnDataChange, WorkerSequence = 1, OnDataChange = new OnDataChangeEvent() },
});
workerClient.Events.Add(new WorkerEvent
{
Event = new MxEvent { Family = MxEventFamily.OnDataChange, WorkerSequence = 2, OnDataChange = new OnDataChangeEvent() },
});
// Constructed but neither worker-attached nor Ready — the premature-attach case.
await using GatewaySession session = CreateSession();
SessionManagerException exception = Assert.Throws<SessionManagerException>(
() => session.AttachInternalEventSubscriber());
Assert.Equal(SessionManagerErrorCode.SessionNotReady, exception.ErrorCode);
// Drive the session to Ready and stream: the failed attach must not have poisoned
// (or even created) the distributor, so a normal subscriber still receives events.
session.AttachWorkerClient(workerClient);
session.MarkReady();
using IEventSubscriberLease lease = session.AttachEventSubscriber(maxSubscribers: 1);
List<MxEvent> received = [];
using CancellationTokenSource readCts = new(TimeSpan.FromSeconds(5));
await foreach (MxEvent mxEvent in lease.Reader.ReadAllAsync(readCts.Token))
{
received.Add(mxEvent);
if (received.Count == 2)
{
break;
}
}
Assert.Equal([1UL, 2UL], received.Select(mxEvent => mxEvent.WorkerSequence).ToArray());
}
private static GatewaySession CreateSession()
{
return new GatewaySession(
sessionId: "session-test-internal-attach",
backendName: "mxaccess",
pipeName: "mxaccess-gateway-1-session-test-internal-attach",
nonce: "nonce",
clientIdentity: "client-1",
ownerKeyId: null,
clientSessionName: "test-session",
clientCorrelationId: "client-correlation-1",
commandTimeout: TimeSpan.FromSeconds(5),
startupTimeout: TimeSpan.FromSeconds(5),
shutdownTimeout: TimeSpan.FromSeconds(5),
leaseDuration: TimeSpan.FromMinutes(30),
openedAt: DateTimeOffset.UtcNow,
eventStreaming: new SessionEventStreaming(
new MxAccessGrpcMapper(),
new EventOptions { QueueCapacity = 8 },
NullLogger<SessionEventDistributor>.Instance,
TimeProvider.System,
new GatewayMetrics()));
}
private static GatewaySession CreateReadySessionWithDetachGrace(
IWorkerClient workerClient,
TimeProvider timeProvider,
@@ -855,6 +930,9 @@ public sealed class GatewaySessionTests
/// <summary>Gets the count of dispose invocations.</summary>
public int DisposeCount { get; private set; }
/// <summary>Events <see cref="ReadEventsAsync"/> yields, in order, before completing. Empty by default.</summary>
public List<WorkerEvent> Events { get; } = [];
/// <inheritdoc />
public Task StartAsync(CancellationToken cancellationToken) => Task.CompletedTask;
@@ -869,7 +947,11 @@ public sealed class GatewaySessionTests
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await Task.CompletedTask.ConfigureAwait(false);
yield break;
foreach (WorkerEvent workerEvent in Events)
{
cancellationToken.ThrowIfCancellationRequested();
yield return workerEvent;
}
}
/// <inheritdoc />
@@ -579,14 +579,6 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
return AsyncEnumerable.Empty<WorkerEvent>();
}
/// <inheritdoc />
public IAsyncEnumerable<MxEvent> ReadAlarmEventsAsync(
string sessionId,
CancellationToken cancellationToken)
{
return AsyncEnumerable.Empty<MxEvent>();
}
/// <inheritdoc />
public Task<SessionCloseResult> CloseSessionAsync(
string sessionId,
@@ -0,0 +1,66 @@
using System.Runtime.CompilerServices;
using System.Threading.Channels;
using ZB.MOM.WW.MxGateway.Contracts.Proto;
using ZB.MOM.WW.MxGateway.Server.Workers;
namespace ZB.MOM.WW.MxGateway.Tests.TestSupport;
/// <summary>
/// Always-<see cref="WorkerClientState.Ready"/> <see cref="IWorkerClient"/> whose event
/// stream is a channel the test writes to. Lets a test attach a real
/// <c>GatewaySession</c> to a worker it drives by hand — the session can be marked Ready,
/// its <c>SessionEventDistributor</c> pump then drains this channel, and the test controls
/// exactly when each event is fanned.
/// </summary>
/// <remarks>
/// Commands are not scripted here: <see cref="InvokeAsync"/> returns an empty reply, because
/// the consumers of this fake route commands through their own <c>ISessionManager</c> double
/// rather than through the worker client. Use a purpose-built worker client instead when a
/// test needs command behavior.
/// </remarks>
/// <param name="sessionId">Session identifier the client reports.</param>
/// <param name="events">Channel whose events <see cref="ReadEventsAsync"/> yields, in order.</param>
public sealed class ChannelWorkerClient(string sessionId, ChannelReader<WorkerEvent> events) : IWorkerClient
{
/// <inheritdoc />
public string SessionId { get; } = sessionId;
/// <inheritdoc />
public int? ProcessId { get; } = 4321;
/// <inheritdoc />
public WorkerClientState State { get; } = WorkerClientState.Ready;
/// <inheritdoc />
public DateTimeOffset LastHeartbeatAt { get; } = DateTimeOffset.UtcNow;
/// <inheritdoc />
public Task StartAsync(CancellationToken cancellationToken) => Task.CompletedTask;
/// <inheritdoc />
public Task<WorkerCommandReply> InvokeAsync(
WorkerCommand command,
TimeSpan timeout,
CancellationToken cancellationToken) => Task.FromResult(new WorkerCommandReply());
/// <inheritdoc />
public async IAsyncEnumerable<WorkerEvent> ReadEventsAsync(
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (WorkerEvent workerEvent in events.ReadAllAsync(cancellationToken).ConfigureAwait(false))
{
yield return workerEvent;
}
}
/// <inheritdoc />
public Task ShutdownAsync(TimeSpan timeout, CancellationToken cancellationToken) => Task.CompletedTask;
/// <inheritdoc />
public void Kill(string reason)
{
}
/// <inheritdoc />
public ValueTask DisposeAsync() => ValueTask.CompletedTask;
}