using System.Diagnostics; using System.Runtime.CompilerServices; using Microsoft.Extensions.Logging; using ZB.MOM.WW.MxGateway.Contracts.Proto; using ZB.MOM.WW.MxGateway.Server.Configuration; using ZB.MOM.WW.MxGateway.Server.Dashboard.Hubs; using ZB.MOM.WW.MxGateway.Server.Grpc; using ZB.MOM.WW.MxGateway.Server.Metrics; using ZB.MOM.WW.MxGateway.Server.Workers; namespace ZB.MOM.WW.MxGateway.Server.Sessions; public sealed class GatewaySession { private readonly object _syncRoot = new(); private readonly SemaphoreSlim _closeLock = new(1, 1); private readonly SessionEventStreaming _eventStreaming; private IWorkerClient? _workerClient; private SessionState _state = SessionState.Creating; private string? _finalFault; private DateTimeOffset _lastClientActivityAt; private DateTimeOffset? _leaseExpiresAt; private bool _closeStarted; private int _activeEventSubscriberCount; private readonly TimeSpan _detachGrace; private readonly TimeSpan _faultedGrace; private readonly TimeSpan _workerReadyWaitTimeout; private DateTimeOffset? _detachedAtUtc; private DateTimeOffset? _faultedAtUtc; // True once at least one external subscriber attached SUCCESSFULLY. Detach-grace's // "last subscriber dropped" stamp (see DetachEventSubscriber) is gated on this so a // FAILED first attach — which still runs the rollback DetachEventSubscriber from the // attach catch path — does not push a never-subscribed session into the grace window. private bool _everHadEventSubscriber; private SessionEventDistributor? _eventDistributor; private bool _eventDistributorStarted; private bool _dashboardMirrorStarted; private IEventSubscriberLease? _dashboardMirrorLease; private Task? _dashboardMirrorTask; private CancellationTokenSource? _dashboardMirrorCts; private readonly Dictionary<(int ServerHandle, int ItemHandle), SessionItemRegistration> _items = []; private readonly ArrayAddressNormalizer? _addressNormalizer; /// /// Initializes a gateway session with session metadata and timeout configuration. /// /// Identifier of the session. /// Name of the backend MXAccess proxy server. /// Name of the named pipe for gateway-worker IPC. /// Security nonce for worker validation. /// Client identity from the authentication context. /// Client-supplied session name. /// Client-supplied correlation identifier. /// Timeout for command invocation. /// Timeout for worker process startup. /// Timeout for worker process shutdown. /// Timestamp when the session opened. /// /// Constructs a session with no owner key ( will be null). /// Authenticated call sites that have a resolved API key identity must use the /// 12-parameter overload and pass the caller's key id explicitly. /// public GatewaySession( string sessionId, string backendName, string pipeName, string nonce, string? clientIdentity, string? clientSessionName, string? clientCorrelationId, TimeSpan commandTimeout, TimeSpan startupTimeout, TimeSpan shutdownTimeout, DateTimeOffset openedAt) : this( sessionId, backendName, pipeName, nonce, clientIdentity, ownerKeyId: null, clientSessionName, clientCorrelationId, commandTimeout, startupTimeout, shutdownTimeout, TimeSpan.FromMinutes(30), openedAt) { } /// /// Initializes a gateway session with session metadata, timeout configuration, and custom lease duration. /// /// Identifier of the session. /// Name of the backend MXAccess proxy server. /// Name of the named pipe for gateway-worker IPC. /// Security nonce for worker validation. /// Client identity from the authentication context. /// API key identifier of the caller that created this session. /// Client-supplied session name. /// Client-supplied correlation identifier. /// Timeout for command invocation. /// Timeout for worker process startup. /// Timeout for worker process shutdown. /// Duration of the session lease. /// Timestamp when the session opened. /// /// Dependencies the session uses to construct and own its /// (the single per-session worker-event pump /// that fans raw mapped s to every subscriber lease). When /// , defaults are used (no replay logger, system clock, a /// fresh mapper, and default ) so unit tests that build a /// session directly still get a working distributor. Production passes the /// DI-resolved dependencies. /// /// /// Retention window kept after the last external (gRPC) event subscriber drops, so a /// client can reconnect. When the window is positive and the active external /// subscriber count falls to zero, the session stays /// and records a detached timestamp; the lease monitor closes it once the window /// elapses with no subscriber having re-attached. (the /// default) disables retention and preserves the original lease-only expiry behavior. /// The clock comes from 's /// so the timer is unit-testable. /// /// /// Bounded time the session will wait, on the command/event hot path, for the worker /// client to reach when the session is already /// but the worker state has transiently diverged /// (e.g. after a heartbeat blip). The wait /// applies only to transient worker states; terminal states /// (// /// /no worker) and a non-Ready session fail /// fast immediately. (the default) disables the wait and /// preserves the original fail-fast behavior byte-for-byte. /// /// /// Rewrites bare array AddItem/AddItem2 addresses to their writable [] /// form using Galaxy metadata at the outbound choke point (and on registration tracking). /// When (legacy unit-construction paths that do not exercise Galaxy /// metadata), addresses pass through unchanged. /// /// /// Grace window kept after the session faults before the lease monitor reaps it. When the /// window is positive the faulted session stays observable via GetSessionStatus for /// that long before it is reclaimed; (the default) makes the /// session reapable on the next sweep. The fault timestamp is stamped in /// using 's clock so the timer /// is unit-testable. /// public GatewaySession( string sessionId, string backendName, string pipeName, string nonce, string? clientIdentity, string? ownerKeyId, string? clientSessionName, string? clientCorrelationId, TimeSpan commandTimeout, TimeSpan startupTimeout, TimeSpan shutdownTimeout, TimeSpan leaseDuration, DateTimeOffset openedAt, SessionEventStreaming? eventStreaming = null, TimeSpan detachGrace = default, TimeSpan workerReadyWaitTimeout = default, ArrayAddressNormalizer? addressNormalizer = null, TimeSpan faultedGrace = default) { if (string.IsNullOrWhiteSpace(sessionId)) { throw new ArgumentException("Session id is required.", nameof(sessionId)); } if (string.IsNullOrWhiteSpace(backendName)) { throw new ArgumentException("Backend name is required.", nameof(backendName)); } if (string.IsNullOrWhiteSpace(pipeName)) { throw new ArgumentException("Pipe name is required.", nameof(pipeName)); } if (string.IsNullOrWhiteSpace(nonce)) { throw new ArgumentException("Nonce is required.", nameof(nonce)); } SessionId = sessionId; BackendName = backendName; PipeName = pipeName; Nonce = nonce; ClientIdentity = clientIdentity; OwnerKeyId = ownerKeyId; ClientSessionName = clientSessionName; ClientCorrelationId = clientCorrelationId; CommandTimeout = commandTimeout; StartupTimeout = startupTimeout; ShutdownTimeout = shutdownTimeout; LeaseDuration = leaseDuration; OpenedAt = openedAt; _lastClientActivityAt = openedAt; _leaseExpiresAt = openedAt + leaseDuration; _eventStreaming = eventStreaming ?? SessionEventStreaming.Default; _detachGrace = detachGrace > TimeSpan.Zero ? detachGrace : TimeSpan.Zero; _faultedGrace = faultedGrace > TimeSpan.Zero ? faultedGrace : TimeSpan.Zero; _workerReadyWaitTimeout = workerReadyWaitTimeout > TimeSpan.Zero ? workerReadyWaitTimeout : TimeSpan.Zero; _addressNormalizer = addressNormalizer; } /// /// Gets the session identifier. /// public string SessionId { get; } /// /// Gets the backend MXAccess proxy server name. /// public string BackendName { get; } /// /// Gets the named pipe name for gateway-worker IPC. /// public string PipeName { get; } /// /// Gets the security nonce for worker validation. /// public string Nonce { get; } /// /// Gets the client identity from the authentication context. /// public string? ClientIdentity { get; } /// /// Gets the API key identifier of the caller that created this session. /// public string? OwnerKeyId { get; } /// /// Gets the client-supplied session name. /// public string? ClientSessionName { get; } /// /// Gets the client-supplied correlation identifier. /// public string? ClientCorrelationId { get; } /// /// Gets the command invocation timeout. /// public TimeSpan CommandTimeout { get; } /// /// Gets the worker process startup timeout. /// public TimeSpan StartupTimeout { get; } /// /// Gets the worker process shutdown timeout. /// public TimeSpan ShutdownTimeout { get; } /// Gets the lease duration for the session. public TimeSpan LeaseDuration { get; } /// /// Gets the timestamp when the session opened. /// public DateTimeOffset OpenedAt { get; } /// /// Gets the worker process identifier, or null if not yet attached. /// public int? WorkerProcessId => _workerClient?.ProcessId; /// /// Gets the attached worker client, or null if not yet attached. /// public IWorkerClient? WorkerClient => _workerClient; /// /// Gets the current session state. /// public SessionState State { get { lock (_syncRoot) { return _state; } } } /// /// Gets the timestamp of the most recent client activity. /// public DateTimeOffset LastClientActivityAt { get { lock (_syncRoot) { return _lastClientActivityAt; } } } /// /// Gets the lease expiration timestamp, or null if no lease is active. /// public DateTimeOffset? LeaseExpiresAt { get { lock (_syncRoot) { return _leaseExpiresAt; } } } /// /// Gets the fault description if the session is faulted, or null. /// public string? FinalFault { get { lock (_syncRoot) { return _finalFault; } } } /// /// Gets the count of active event stream subscribers. /// public int ActiveEventSubscriberCount { get { lock (_syncRoot) { return _activeEventSubscriberCount; } } } /// /// Gets the UTC timestamp at which the session entered its detach-grace retention /// window (the last external event subscriber dropped while a positive /// detach-grace was configured), or when the session is not /// currently within a detach-grace window. Re-attaching an external subscriber clears /// this. Always when detach-grace is disabled /// (DetachGraceSeconds == 0). /// public DateTimeOffset? DetachedAtUtc { get { lock (_syncRoot) { return _detachedAtUtc; } } } /// /// Attaches the worker client for this session. /// /// Worker client to attach. public void AttachWorkerClient(IWorkerClient workerClient) { ArgumentNullException.ThrowIfNull(workerClient); lock (_syncRoot) { _workerClient = workerClient; } } /// /// Transitions the session to a new state with constraints for terminal states. /// /// Next session state to transition to. /// /// is terminal. /// only allows a transition to . /// only allows a transition to /// (or ) — once /// has started, no late lifecycle callback can revive the /// session by walking it back to or any earlier /// state. Both close-related writes (Closing and Closed) go through /// _syncRoot just like every other state read/write, closing the split-lock /// race. /// public void TransitionTo(SessionState nextState) { lock (_syncRoot) { if (_state is SessionState.Closed) { return; } if (_state is SessionState.Faulted && nextState is not SessionState.Closed) { return; } if (_state is SessionState.Closing && nextState is not SessionState.Closed && nextState is not SessionState.Faulted) { return; } _state = nextState; } } /// /// Transitions the session to the Ready state. /// /// /// On becoming Ready the session starts its internal dashboard mirror when a /// dashboard broadcaster was supplied. The mirror registers an internal subscriber on /// the distributor and starts the pump before any gRPC client attaches, so the /// dashboard EventsHub receives session events even with no gRPC subscriber streaming — /// fixing the "dark feed" where the dashboard only saw events while a gRPC client was /// actively streaming. Registering the internal subscriber BEFORE /// also avoids the hazard where /// starting the pump at Ready with zero subscribers drained a fast-completing worker /// stream into nothing and left a later subscriber hanging: there is now always a /// subscriber (the dashboard one) registered before the pump starts. /// public void MarkReady() { TransitionTo(SessionState.Ready); StartDashboardMirror(); } // Constructs and starts the distributor exactly once, registering the subscriber under // the same start so no event the pump fans can be missed between start and register. // Started lazily on the FIRST AttachEventSubscriber rather than at MarkReady: today the // worker event stream is only drained when a client begins streaming, so deferring the // single drain to first-attach preserves that "events start flowing on subscribe" // behavior and avoids draining a fast-completing source into the void before any // subscriber exists. The source factory mirrors the mapping/ordering/start that // EventStreamService.ProduceEventsAsync previously used: it drains the worker event // stream in source order and maps each WorkerEvent to the public MxEvent with the same // mapper, with no skip/filter — per-RPC filtering (e.g. AfterWorkerSequence) stays at the // subscriber boundary in EventStreamService. Returns a registered lease atomically with // the start so the very first subscriber sees the stream from its beginning. private IEventSubscriberLease StartDistributorAndRegister() { SessionEventDistributor distributor = EnsureDistributorCreated(out bool startNow); // Register BEFORE starting the pump so a subscriber is present when the pump begins // draining — no event is fanned to an empty subscriber set and then missed by this // first subscriber. StartAsync only schedules the pump task; it never blocks. IEventSubscriberLease lease = distributor.Register(); StartPumpIfRequested(distributor, startNow); return lease; } // Reconnect/resume variant of StartDistributorAndRegister. Snapshots the replay // ring for events newer than afterSequence AND registers the live subscriber atomically // under the distributor's replay lock, so the replay→live handoff has no gap and no // duplicate (see SessionEventDistributor.RegisterWithReplay). The pump is started after // registration, exactly as the fresh-attach path, so the very first subscriber on a // freshly-Ready session still sees the stream from its beginning. private IEventSubscriberLease StartDistributorAndRegisterWithReplay( ulong afterSequence, out IReadOnlyList replayedEvents, out bool gap, out ulong oldestAvailableSequence, out ulong liveResumeSequence) { SessionEventDistributor distributor = EnsureDistributorCreated(out bool startNow); IEventSubscriberLease lease = distributor.RegisterWithReplay( afterSequence, out replayedEvents, out gap, out oldestAvailableSequence, out liveResumeSequence); StartPumpIfRequested(distributor, startNow); return lease; } // Constructs the distributor exactly once and reports whether THIS caller is the one // that should start the pump (i.e. it observed the unstarted state and claimed the // start). Both the construction and the started-flag flip happen under _syncRoot so two // concurrent callers (e.g. MarkReady's dashboard mirror and a racing first // AttachEventSubscriber) agree on a single distributor and a single start. private SessionEventDistributor EnsureDistributorCreated(out bool startNow) { lock (_syncRoot) { if (_eventDistributor is null) { EventOptions eventOptions = _eventStreaming.EventOptions; _eventDistributor = new SessionEventDistributor( SessionId, MapWorkerEventsAsync, eventOptions.QueueCapacity, eventOptions.ReplayBufferCapacity, eventOptions.ReplayRetentionSeconds, _eventStreaming.DistributorLogger, _eventStreaming.TimeProvider, CreateOverflowHandler(eventOptions.BackpressurePolicy), singleSubscriberMode: !_eventStreaming.AllowMultipleEventSubscribers); } startNow = false; if (!_eventDistributorStarted) { _eventDistributorStarted = true; startNow = true; } return _eventDistributor; } } /// /// Registers a gateway-owned internal (non-counted) distributor subscriber and /// returns its lease. The lease's yields the /// same mapped s the single distributor pump fans to every /// subscriber; disposing the lease unregisters it. /// /// /// Used by the central alarm monitor so it consumes events through the one distributor /// pump instead of opening a second raw drain of the single worker event channel (which /// would split events between the two readers). Mirrors the dashboard-mirror lease: /// isInternal: true keeps this subscriber out of the /// MaxEventSubscribersPerSession 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. /// /// Gated on readiness exactly like : attaching /// before the session and its worker are Ready throws /// with /// . /// /// /// The internal subscriber's lease; dispose it to unregister. /// /// The session or its worker client is not Ready. /// 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. SessionEventDistributor distributor = EnsureDistributorCreated(out bool startNow); IEventSubscriberLease lease = distributor.Register(isInternal: true); StartPumpIfRequested(distributor, startNow); return lease; } private static void StartPumpIfRequested(SessionEventDistributor distributor, bool startNow) { if (!startNow) { return; } // StartAsync only schedules the pump via Task.Run and returns a completed task; // it does not perform any async I/O itself. The sync-over-async call here is // therefore safe and will not deadlock. Do not make StartAsync truly async // (i.e., await real I/O before returning) without also changing this call site. distributor.StartAsync(CancellationToken.None).GetAwaiter().GetResult(); } // Registers the gateway-owned internal dashboard subscriber on the distributor and starts // a background loop that mirrors every fanned event to the dashboard broadcaster. Called // once when the session becomes Ready (idempotent). The internal subscriber is registered // BEFORE the pump starts (see StartDistributorAndRegister / EnsureDistributorCreated), so // a subscriber is always present at pump start — the dashboard receives events with no // gRPC subscriber attached, and the "zero-subscriber drain into the void" hang // cannot occur. No-op when no dashboard broadcaster was supplied (unit tests). // // Race-safety (Issue 1): _dashboardMirrorLease and _dashboardMirrorTask are published // atomically under a SINGLE second lock section, and DisposeAsync reads/nulls them under // that same lock. After EnsureDistributorCreated/Register/StartPump (all outside _syncRoot // to avoid lock inversion with the distributor's own lifecycle lock), we re-enter // _syncRoot and check for concurrent disposal. If the session is already Closing/Closed/ // Faulted at that point, we dispose the just-created lease immediately and do NOT start // the mirror task, so nothing is orphaned. private void StartDashboardMirror() { IDashboardEventBroadcaster? broadcaster = _eventStreaming.DashboardBroadcaster; if (broadcaster is null) { return; } CancellationToken loopToken; lock (_syncRoot) { if (_dashboardMirrorStarted || _state is SessionState.Closing or SessionState.Closed or SessionState.Faulted) { return; } _dashboardMirrorStarted = true; _dashboardMirrorCts = new CancellationTokenSource(); loopToken = _dashboardMirrorCts.Token; } // Create the distributor (claiming the start if we are first) and register the // internal subscriber BEFORE starting the pump. isInternal: true keeps the dashboard // subscriber out of the single-subscriber overflow accounting, so a slow/broken // dashboard mirror only disconnects itself and never faults the session. // These three calls are OUTSIDE _syncRoot to avoid holding it across // EnsureDistributorCreated's own lock and StartAsync's Task.Run. SessionEventDistributor distributor = EnsureDistributorCreated(out bool startNow); IEventSubscriberLease lease = distributor.Register(isInternal: true); StartPumpIfRequested(distributor, startNow); // Publish BOTH the lease and the task atomically under one lock section so // DisposeAsync always sees them in a consistent state: either both are set or // both are null. If the session already started disposal before we got here, // dispose the lease immediately instead of orphaning it. lock (_syncRoot) { if (_state is SessionState.Closing or SessionState.Closed or SessionState.Faulted) { // Disposal already ran (or is in progress) — discard the just-created // lease now so it is not orphaned. Do NOT launch the mirror task. lease.Dispose(); return; } _dashboardMirrorLease = lease; _dashboardMirrorTask = Task.Run( () => RunDashboardMirrorAsync(broadcaster, lease, loopToken), CancellationToken.None); } } // Reads the internal dashboard subscriber's channel and publishes each RAW fanned event // to the dashboard broadcaster. The dashboard is a first-class distributor subscriber, // so it sees the session's full raw event activity — NOT the per-gRPC-subscriber // AfterWorkerSequence filtering that EventStreamService applies at its own boundary. This // is intentional: the dashboard is a separate LDAP-authenticated monitoring view (per- // session dashboard ACL is a separate concern). Publish is best-effort / never-throw, so // a slow or broken dashboard cannot fault the session or stall the pump; the bounded // internal subscriber channel only disconnects THIS mirror on overflow, leaving the // session and other subscribers untouched. private async Task RunDashboardMirrorAsync( IDashboardEventBroadcaster broadcaster, IEventSubscriberLease lease, CancellationToken cancellationToken) { try { await foreach (MxEvent mxEvent in lease.Reader .ReadAllAsync(cancellationToken) .ConfigureAwait(false)) { try { broadcaster.Publish(SessionId, mxEvent); } catch (Exception exception) { // Publish is documented never-throw, but enforce it here too so a future // implementation cannot fault the mirror loop. Logs identifiers only. _eventStreaming.DistributorLogger.LogDebug( exception, "Dashboard event mirror threw for session {SessionId}; continuing.", SessionId); } } } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { // Teardown path: the session is shutting down the mirror. } catch (SessionManagerException) { // The internal subscriber's channel overflowed and the distributor disconnected // it with a terminal overflow fault. That disconnects only the dashboard mirror; // the session, pump, and any gRPC subscriber are unaffected. Stop mirroring. } catch (Exception exception) { // Source-fault completion (worker event stream terminated abnormally) surfaces // here. The session's own fault handling runs via the gRPC path / lifecycle; the // mirror just stops. Logs identifiers only. _eventStreaming.DistributorLogger.LogDebug( exception, "Dashboard event mirror loop ended for session {SessionId}.", SessionId); } } // Builds the per-subscriber backpressure handler the distributor invokes when a // subscriber's bounded channel overflows. The distributor always disconnects the // offending subscriber with an EventQueueOverflow fault; this handler adds the // observable side effects, preserving exactly what the pre-epic per-RPC overflow path // emitted: // - always record the queue-overflow metric, labeled by subscriber kind; // - FailFast in the legacy single-subscriber case (isOnlySubscriber): fault the whole // session and record the fault metric, matching back-compat behavior; // - FailFast with multiple subscribers, or DisconnectSubscriber in any case: do NOT // fault the session — the distributor's disconnect of the one slow subscriber is the // whole remedy, so other subscribers and the pump are unaffected. Multi-subscriber // FailFast deliberately degrades to a disconnect because faulting a shared session on // one slow consumer would punish healthy subscribers. // The delegate now carries isInternal directly (Issue 4), so the metric label is chosen // without any heuristic: "dashboard-mirror" for internal, "grpc-event-stream" for external. private SubscriberOverflowHandler CreateOverflowHandler(EventBackpressurePolicy policy) { GatewayMetrics metrics = _eventStreaming.Metrics; string sessionId = SessionId; return (isOnlySubscriber, isInternal) => { // Label the overflow metric by subscriber kind. The distributor passes isInternal // directly, so no heuristic is needed to distinguish an internal overflow (the // gateway-owned dashboard mirror) from an external one (a gRPC streaming client). string label = isInternal ? "dashboard-mirror" : "grpc-event-stream"; metrics.QueueOverflow(label); if (policy == EventBackpressurePolicy.FailFast && isOnlySubscriber) { MarkFaulted($"Session {sessionId} event stream queue overflowed."); metrics.Fault(SessionManagerErrorCode.EventQueueOverflow.ToString()); } }; } // The distributor's single event source. Drains the worker event stream once (the // distributor guarantees a single consumer) and maps each frame to the public MxEvent, // preserving worker order. Mirrors the former ProduceEventsAsync mapping exactly. // // This is the session's only reader of the worker event channel: every gateway consumer — // gRPC subscribers, the dashboard mirror, the alarm monitor — attaches to the distributor // this source feeds. WorkerClient.ReadEventsAsync single-reader-claims that channel and // throws on a second consumer, so any future path that drains it directly fails loudly // rather than splitting events. private async IAsyncEnumerable MapWorkerEventsAsync( [EnumeratorCancellation] CancellationToken cancellationToken) { MxAccessGrpcMapper mapper = _eventStreaming.Mapper; IWorkerClient workerClient = await GetReadyWorkerClientAsync(cancellationToken).ConfigureAwait(false); TouchClientActivity(_eventStreaming.TimeProvider.GetUtcNow()); await foreach (WorkerEvent workerEvent in workerClient .ReadEventsAsync(cancellationToken) .WithCancellation(cancellationToken) .ConfigureAwait(false)) { yield return mapper.MapEvent(workerEvent); } } /// /// Transitions the session to the Faulted state with a fault description. /// /// Reason for the fault. public void MarkFaulted(string reason) { lock (_syncRoot) { if (_state is SessionState.Closed) { return; } _finalFault = reason; _state = SessionState.Faulted; // Stamp the fault time once, on the first fault, so the sweeper can apply // FaultedGraceSeconds. A subsequent MarkFaulted (already-faulted session) keeps the // original timestamp so the grace window is measured from the first fault. _faultedAtUtc ??= _eventStreaming.TimeProvider.GetUtcNow(); } } /// /// Updates the timestamp of the most recent client activity. /// /// Timestamp of the client activity. public void TouchClientActivity(DateTimeOffset activityAt) { lock (_syncRoot) { _lastClientActivityAt = activityAt; _leaseExpiresAt = activityAt + LeaseDuration; } } /// /// Extends the session lease to the specified expiration time. /// /// Timestamp when the lease expires. public void ExtendLease(DateTimeOffset leaseExpiresAt) { lock (_syncRoot) { _leaseExpiresAt = leaseExpiresAt; } } /// /// Determines whether the session lease has expired. /// /// Current timestamp for comparison. /// if the lease has expired with no active event subscriber; otherwise . public bool IsLeaseExpired(DateTimeOffset now) { lock (_syncRoot) { return _activeEventSubscriberCount == 0 && _leaseExpiresAt is not null && _leaseExpiresAt <= now; } } /// /// Determines whether the session's detach-grace retention window has elapsed: the /// session entered detach-grace (its last external event subscriber dropped while a /// positive detach-grace was configured) and has had no external subscriber re-attach /// for longer than the configured detach-grace. The lease monitor closes such a /// session exactly as it closes an expired lease. Always returns /// when detach-grace is disabled or when an external subscriber is attached (the /// detached timestamp is cleared on re-attach, so an attached session is never within a /// window). /// /// Current timestamp for comparison. /// if the detach-grace window has elapsed with no re-attached subscriber; otherwise . public bool IsDetachGraceExpired(DateTimeOffset now) { lock (_syncRoot) { return _detachGrace > TimeSpan.Zero && _activeEventSubscriberCount == 0 && _detachedAtUtc is not null && now - _detachedAtUtc.Value >= _detachGrace; } } /// /// Determines whether a faulted session is now eligible for reaping by the lease monitor. /// A faulted session is permanently unusable (every command fails the readiness check), /// so the sweeper closes it exactly as it closes an expired lease — but no sooner than the /// configured FaultedGraceSeconds after the fault, so a monitoring client can still /// observe the fault before the slot is reclaimed. Always returns /// for a non-faulted session. /// /// Current timestamp for comparison. /// if the session is faulted and past its fault-grace window; otherwise . public bool IsFaultedReapable(DateTimeOffset now) { lock (_syncRoot) { return IsFaultedReapableCore(now); } } /// /// Attaches an event subscriber and returns a lease whose /// reads the fanned public /// s for this subscriber. The returned lease, when disposed, /// unregisters the distributor subscriber AND decrements the active-subscriber count. /// /// /// Maximum concurrent external subscribers in multi-subscriber mode /// (MxGateway:Sessions:MaxEventSubscribersPerSession). Ignored when the /// session is in single-subscriber mode (AllowMultipleEventSubscribers == false); /// the effective cap is then 1. The gateway-owned internal dashboard subscriber is /// registered directly on the distributor and is NOT counted here, so it never /// consumes cap budget. /// /// /// The subscriber mode is derived internally from /// — the same source /// the uses to gate its FailFast decision — so /// the cap-enforcement mode and the distributor's singleSubscriberMode field /// cannot diverge. The count-check-and-increment runs atomically under /// _syncRoot, so two concurrent attaches racing toward the cap can never both /// succeed past it. On distributor-register failure the count is rolled back (see the /// catch below). /// /// A lease that reads the fanned public events for this subscriber. public IEventSubscriberLease AttachEventSubscriber(int maxSubscribers) { // Derive the mode from the same source the distributor uses so the two can never // diverge. Effective cap: 1 in single-subscriber mode, otherwise the configured // maximum (clamped to at least 1 so a misconfigured non-positive value can never // deadlock attaches in multi-subscriber mode). bool allowMultipleSubscribers = _eventStreaming.AllowMultipleEventSubscribers; int effectiveCap = allowMultipleSubscribers ? Math.Max(1, maxSubscribers) : 1; 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}."); } if (_activeEventSubscriberCount >= effectiveCap) { throw allowMultipleSubscribers ? new SessionManagerException( SessionManagerErrorCode.EventSubscriberLimitReached, $"Session {SessionId} has reached its maximum of {effectiveCap} concurrent event stream subscribers.") : new SessionManagerException( SessionManagerErrorCode.EventSubscriberAlreadyActive, $"Session {SessionId} already has an active event stream subscriber."); } _activeEventSubscriberCount++; // An external subscriber (re)attached: cancel any in-flight detach-grace window so // the lease monitor no longer treats this session as eligible for grace-expiry // close. This is the reattach→grace-cancel transition; it races the sweeper's // IsDetachGraceExpired read, and both run under _syncRoot so they serialize. _detachedAtUtc = null; } // Construct/start the distributor and register this subscriber. Done outside the // guard lock (StartDistributorAndRegister takes _syncRoot itself for construction). // On any failure roll back the count we just took so the guard stays consistent. try { IEventSubscriberLease distributorLease = StartDistributorAndRegister(); MarkEventSubscriberAttached(); return new EventSubscriberLease(this, distributorLease); } catch { DetachEventSubscriber(); throw; } } /// /// Reconnect/resume variant of . Attaches /// an event subscriber AND atomically snapshots the session replay ring for events newer /// than , so a resuming client can replay what it missed /// before live delivery resumes — with no gap and no duplicate across the handoff. /// /// See . /// /// The last worker sequence the resuming client already observed. Replay returns events /// strictly newer than this; the caller must filter the live channel to events strictly /// newer than . /// /// /// The lease plus the replay batch, gap flag, and resume watermarks. See /// for the no-gap/no-duplicate /// guarantee. /// public EventSubscriberReplayAttachment AttachEventSubscriberWithReplay(int maxSubscribers, ulong afterSequence) { bool allowMultipleSubscribers = _eventStreaming.AllowMultipleEventSubscribers; int effectiveCap = allowMultipleSubscribers ? Math.Max(1, maxSubscribers) : 1; 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}."); } if (_activeEventSubscriberCount >= effectiveCap) { throw allowMultipleSubscribers ? new SessionManagerException( SessionManagerErrorCode.EventSubscriberLimitReached, $"Session {SessionId} has reached its maximum of {effectiveCap} concurrent event stream subscribers.") : new SessionManagerException( SessionManagerErrorCode.EventSubscriberAlreadyActive, $"Session {SessionId} already has an active event stream subscriber."); } _activeEventSubscriberCount++; _detachedAtUtc = null; } try { IEventSubscriberLease distributorLease = StartDistributorAndRegisterWithReplay( afterSequence, out IReadOnlyList replayedEvents, out bool gap, out ulong oldestAvailableSequence, out ulong liveResumeSequence); MarkEventSubscriberAttached(); return new EventSubscriberReplayAttachment( new EventSubscriberLease(this, distributorLease), replayedEvents, gap, oldestAvailableSequence, liveResumeSequence); } catch { DetachEventSubscriber(); throw; } } // Records that an external subscriber attached successfully. Gates the detach-grace // "last subscriber dropped" stamp so a FAILED first attach (which still rolls back via // DetachEventSubscriber) never pushes a never-subscribed session into grace. private void MarkEventSubscriberAttached() { lock (_syncRoot) { _everHadEventSubscriber = true; } } /// /// Invokes a worker command synchronously and returns the reply. /// /// Worker command to invoke. /// Token to cancel the asynchronous operation. /// The worker's reply to the command. public async Task InvokeAsync( WorkerCommand command, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(command); if (command.Command is not null) { NormalizeOutboundCommand(command.Command); } IWorkerClient workerClient = await GetReadyWorkerClientAsync(cancellationToken).ConfigureAwait(false); TouchClientActivity(_eventStreaming.TimeProvider.GetUtcNow()); return await workerClient.InvokeAsync(command, CommandTimeout, cancellationToken).ConfigureAwait(false); } // Single outbound choke point for the two array-write ergonomics shims: // 1. AddItem/AddItem2 array addresses gain the writable "[]" suffix when Galaxy metadata // reports them as arrays, so the worker registers a write-capable handle. The mutation // lands on the same MxCommand instance forwarded to the worker. // 2. Sparse array write values are expanded to whole-array values, because MXAccess has no // partial-array write primitive — the worker only ever sees a full MxArray. // SparseArrayExpander.Expand throws RpcException(InvalidArgument) for an invalid sparse payload; // that propagates out of InvokeAsync as the desired client-facing error and is deliberately not // caught here. private void NormalizeOutboundCommand(MxCommand command) { switch (command.PayloadCase) { case MxCommand.PayloadOneofCase.AddItem: command.AddItem.ItemDefinition = NormalizeAddress(command.AddItem.ItemDefinition); break; case MxCommand.PayloadOneofCase.AddItem2: command.AddItem2.ItemDefinition = NormalizeAddress(command.AddItem2.ItemDefinition); break; case MxCommand.PayloadOneofCase.AddBufferedItem: command.AddBufferedItem.ItemDefinition = NormalizeAddress(command.AddBufferedItem.ItemDefinition); break; case MxCommand.PayloadOneofCase.AddItemBulk: // Normalize each bare array address in place so the worker binds a write-capable handle // for every array tag in the batch (the same IsArray-gated rewrite the single-add path // applies). Scalar addresses pass through unchanged. for (int i = 0; i < command.AddItemBulk.TagAddresses.Count; i++) { command.AddItemBulk.TagAddresses[i] = NormalizeAddress(command.AddItemBulk.TagAddresses[i]); } break; case MxCommand.PayloadOneofCase.Write: ExpandValue(command.Write.Value); break; case MxCommand.PayloadOneofCase.WriteSecured: ExpandValue(command.WriteSecured.Value); break; case MxCommand.PayloadOneofCase.Write2: ExpandValue(command.Write2.Value); break; case MxCommand.PayloadOneofCase.WriteSecured2: ExpandValue(command.WriteSecured2.Value); break; case MxCommand.PayloadOneofCase.WriteBulk: foreach (WriteBulkEntry entry in command.WriteBulk.Entries) { ExpandValue(entry.Value); } break; case MxCommand.PayloadOneofCase.Write2Bulk: foreach (Write2BulkEntry entry in command.Write2Bulk.Entries) { ExpandValue(entry.Value); } break; case MxCommand.PayloadOneofCase.WriteSecuredBulk: foreach (WriteSecuredBulkEntry entry in command.WriteSecuredBulk.Entries) { ExpandValue(entry.Value); } break; case MxCommand.PayloadOneofCase.WriteSecured2Bulk: foreach (WriteSecured2BulkEntry entry in command.WriteSecured2Bulk.Entries) { ExpandValue(entry.Value); } break; } } // Best-effort array-suffix rewrite; the normalizer is null in legacy unit-construction paths // that do not exercise Galaxy metadata, in which case the address passes through unchanged. private string NormalizeAddress(string address) => _addressNormalizer?.Normalize(address) ?? address; // MXAccess writes replace the whole array; expand a sparse value in place so the worker only // ever receives a whole-array MxValue. No-op for null or non-sparse values. The configured // MxGateway:Events:MaxSparseArrayLength cap is enforced before the full array is allocated. private void ExpandValue(MxValue? value) { if (value is not null) { SparseArrayExpander.Expand(value, _eventStreaming.EventOptions.MaxSparseArrayLength); } } /// Gets the item registration for a server and item handle pair. /// The MXAccess server handle. /// The MXAccess item handle. /// The item registration if found. /// if a registration was found for the handle pair; otherwise . public bool TryGetItemRegistration( int serverHandle, int itemHandle, out SessionItemRegistration registration) { lock (_syncRoot) { return _items.TryGetValue((serverHandle, itemHandle), out registration!); } } /// Tracks item registrations from a command reply. /// The executed command. /// The command reply. public void TrackCommandReply( MxCommand command, MxCommandReply reply) { if (reply.ProtocolStatus?.Code is not ProtocolStatusCode.Ok) { return; } lock (_syncRoot) { switch (command.Kind) { // The public reply is tracked from the pre-mapping MxCommand instance, which is a // separate copy from the one mutated at the InvokeAsync choke point (the gRPC mapper // deep-clones before forwarding). Re-apply the array-suffix normalization here so the // registration's TagAddress matches the address the worker actually registered. // Normalize is idempotent for an already-suffixed address. case MxCommandKind.AddItem when reply.AddItem is not null: TrackItem(command.AddItem.ServerHandle, reply.AddItem.ItemHandle, NormalizeAddress(command.AddItem.ItemDefinition)); break; case MxCommandKind.AddItem2 when reply.AddItem2 is not null: TrackItem(command.AddItem2.ServerHandle, reply.AddItem2.ItemHandle, NormalizeAddress(command.AddItem2.ItemDefinition)); break; case MxCommandKind.AddBufferedItem when reply.AddBufferedItem is not null: // The reply carries no address, so tracking keys off the command's ItemDefinition; // re-apply the array-suffix normalization (the tracking copy is a separate, un-mutated // instance from the one forwarded at the InvokeAsync choke point) so the registration // matches the write-capable handle the worker bound. TrackItem(command.AddBufferedItem.ServerHandle, reply.AddBufferedItem.ItemHandle, NormalizeAddress(command.AddBufferedItem.ItemDefinition)); break; case MxCommandKind.AddItemBulk when reply.AddItemBulk is not null: // The worker echoes back the (already-normalized) address it bound in each // SubscribeResult.TagAddress, so TrackBulkItems stores the suffixed array address // without re-normalizing here. TrackBulkItems(reply.AddItemBulk); break; case MxCommandKind.SubscribeBulk when reply.SubscribeBulk is not null: TrackBulkItems(reply.SubscribeBulk); break; case MxCommandKind.RemoveItem: _items.Remove((command.RemoveItem.ServerHandle, command.RemoveItem.ItemHandle)); break; case MxCommandKind.RemoveItemBulk: RemoveItems(command.RemoveItemBulk.ServerHandle, command.RemoveItemBulk.ItemHandles); break; case MxCommandKind.UnsubscribeBulk: RemoveItems(command.UnsubscribeBulk.ServerHandle, command.UnsubscribeBulk.ItemHandles); break; } } } /// /// Executes a bulk add-item command for the specified server and tag addresses. /// /// Server handle returned by the worker. /// Tag addresses to add. /// Token to cancel the asynchronous operation. /// The per-address subscribe results. public Task> AddItemBulkAsync( int serverHandle, IReadOnlyList tagAddresses, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(tagAddresses); AddItemBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.TagAddresses.Add(tagAddresses); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.AddItemBulk, AddItemBulk = bulkCommand, }, reply => reply.AddItemBulk, cancellationToken); } /// /// Executes a bulk advise-item command for the specified server and item handles. /// /// Server handle returned by the worker. /// Item handles to advise. /// Token to cancel the asynchronous operation. /// The per-handle subscribe results. public Task> AdviseItemBulkAsync( int serverHandle, IReadOnlyList itemHandles, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(itemHandles); AdviseItemBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.ItemHandles.Add(itemHandles); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.AdviseItemBulk, AdviseItemBulk = bulkCommand, }, reply => reply.AdviseItemBulk, cancellationToken); } /// /// Executes a bulk remove-item command for the specified server and item handles. /// /// Server handle returned by the worker. /// Item handles to remove. /// Token to cancel the asynchronous operation. /// The per-handle subscribe results. public Task> RemoveItemBulkAsync( int serverHandle, IReadOnlyList itemHandles, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(itemHandles); RemoveItemBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.ItemHandles.Add(itemHandles); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.RemoveItemBulk, RemoveItemBulk = bulkCommand, }, reply => reply.RemoveItemBulk, cancellationToken); } /// /// Executes a bulk un-advise-item command for the specified server and item handles. /// /// Server handle returned by the worker. /// Item handles to un-advise. /// Token to cancel the asynchronous operation. /// The per-handle subscribe results. public Task> UnAdviseItemBulkAsync( int serverHandle, IReadOnlyList itemHandles, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(itemHandles); UnAdviseItemBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.ItemHandles.Add(itemHandles); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.UnAdviseItemBulk, UnAdviseItemBulk = bulkCommand, }, reply => reply.UnAdviseItemBulk, cancellationToken); } /// /// Executes a bulk subscribe command for the specified server and tag addresses. /// /// Server handle returned by the worker. /// Tag addresses to subscribe to. /// Token to cancel the asynchronous operation. /// The per-address subscribe results. public Task> SubscribeBulkAsync( int serverHandle, IReadOnlyList tagAddresses, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(tagAddresses); SubscribeBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.TagAddresses.Add(tagAddresses); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.SubscribeBulk, SubscribeBulk = bulkCommand, }, reply => reply.SubscribeBulk, cancellationToken); } /// /// Executes a bulk unsubscribe command for the specified server and item handles. /// /// Server handle returned by the worker. /// Item handles to unsubscribe from. /// Token to cancel the asynchronous operation. /// The per-handle subscribe results. public Task> UnsubscribeBulkAsync( int serverHandle, IReadOnlyList itemHandles, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(itemHandles); UnsubscribeBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.ItemHandles.Add(itemHandles); return InvokeBulkAsync( new MxCommand { Kind = MxCommandKind.UnsubscribeBulk, UnsubscribeBulk = bulkCommand, }, reply => reply.UnsubscribeBulk, cancellationToken); } /// Executes a bulk Write command for the specified server and per-item entries. /// Server handle returned by the worker. /// Write entries to execute. /// Token to cancel the asynchronous operation. /// The per-entry write results. public Task> WriteBulkAsync( int serverHandle, IReadOnlyList entries, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(entries); WriteBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.Entries.Add(entries); return InvokeBulkWriteAsync( new MxCommand { Kind = MxCommandKind.WriteBulk, WriteBulk = bulkCommand, }, reply => reply.WriteBulk, cancellationToken); } /// Executes a bulk Write2 (timestamped) command. /// Server handle returned by the worker. /// Write entries to execute. /// Token to cancel the asynchronous operation. /// The per-entry write results. public Task> Write2BulkAsync( int serverHandle, IReadOnlyList entries, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(entries); Write2BulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.Entries.Add(entries); return InvokeBulkWriteAsync( new MxCommand { Kind = MxCommandKind.Write2Bulk, Write2Bulk = bulkCommand, }, reply => reply.Write2Bulk, cancellationToken); } /// Executes a bulk WriteSecured command. /// Server handle returned by the worker. /// Write entries to execute. /// Token to cancel the asynchronous operation. /// The per-entry write results. public Task> WriteSecuredBulkAsync( int serverHandle, IReadOnlyList entries, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(entries); WriteSecuredBulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.Entries.Add(entries); return InvokeBulkWriteAsync( new MxCommand { Kind = MxCommandKind.WriteSecuredBulk, WriteSecuredBulk = bulkCommand, }, reply => reply.WriteSecuredBulk, cancellationToken); } /// Executes a bulk WriteSecured2 command. /// Server handle returned by the worker. /// Write entries to execute. /// Token to cancel the asynchronous operation. /// The per-entry write results. public Task> WriteSecured2BulkAsync( int serverHandle, IReadOnlyList entries, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(entries); WriteSecured2BulkCommand bulkCommand = new() { ServerHandle = serverHandle }; bulkCommand.Entries.Add(entries); return InvokeBulkWriteAsync( new MxCommand { Kind = MxCommandKind.WriteSecured2Bulk, WriteSecured2Bulk = bulkCommand, }, reply => reply.WriteSecured2Bulk, cancellationToken); } /// /// Executes a bulk Read command — see ReadBulkCommand's doc /// comment in the .proto for the cached-vs-snapshot semantics. /// /// Server handle returned by the worker. /// Tag addresses to read. /// Timeout for the read operation. /// Token to cancel the asynchronous operation. /// The per-address read results. public Task> ReadBulkAsync( int serverHandle, IReadOnlyList tagAddresses, TimeSpan timeout, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(tagAddresses); ReadBulkCommand bulkCommand = new() { ServerHandle = serverHandle, TimeoutMs = timeout <= TimeSpan.Zero ? 0u : (uint)Math.Min(timeout.TotalMilliseconds, uint.MaxValue), }; bulkCommand.TagAddresses.Add(tagAddresses); return InvokeBulkReadAsync( new MxCommand { Kind = MxCommandKind.ReadBulk, ReadBulk = bulkCommand, }, reply => reply.ReadBulk, cancellationToken); } /// /// Closes the session and shuts down the worker process. /// /// Reason for closing the session. /// Token to cancel the asynchronous operation. /// /// Concurrent close attempts are serialized by _closeLock so only one close /// runs at a time, but every read/write of _state still passes through /// _syncRoot (via and ) — /// the close path therefore obeys the same lock discipline as /// / and a concurrent /// TransitionTo(Ready) cannot race past a Closing write. /// /// The outcome of the close operation. public async Task CloseAsync( string reason, CancellationToken cancellationToken) { await _closeLock.WaitAsync(cancellationToken).ConfigureAwait(false); try { try { if (!TryBeginClose(out bool alreadyClosing)) { return new SessionCloseResult(SessionId, SessionState.Closed, AlreadyClosed: true); } if (_workerClient is not null) { try { await _workerClient.ShutdownAsync(ShutdownTimeout, cancellationToken).ConfigureAwait(false); } catch (Exception exception) { try { _workerClient.Kill(reason); } catch (Exception killException) { throw new SessionCloseStartedException( $"Session {SessionId} close failed after worker shutdown started.", new AggregateException(exception, killException)); } throw; } } MarkClosed(); return new SessionCloseResult(SessionId, SessionState.Closed, alreadyClosing); } catch (Exception exception) when (exception is not SessionCloseStartedException) { throw new SessionCloseStartedException( $"Session {SessionId} close failed after the close lock was acquired.", exception); } } finally { _closeLock.Release(); } } // Returns false when the session is already Closed (caller short-circuits with // AlreadyClosed: true). Otherwise sets _state = Closing under _syncRoot so a // concurrent TransitionTo(Ready) — which only refuses to overwrite Closed/Faulted // — cannot flip the session back to Ready after close started. The `alreadyClosing` // out parameter mirrors the previous `_closeStarted` check so the surface contract // (a second concurrent close returns AlreadyClosed: alreadyClosing) is preserved. private bool TryBeginClose(out bool alreadyClosing) { lock (_syncRoot) { if (_state is SessionState.Closed) { alreadyClosing = _closeStarted; return false; } alreadyClosing = _closeStarted; _closeStarted = true; _state = SessionState.Closing; return true; } } /// /// Atomically re-verifies that the session is still eligible for sweep-initiated close /// (lease expired OR detach-grace expired, with no active external subscriber) and, if so, /// transitions to Closing in a single lock acquisition. /// /// Current timestamp used for expiry re-check. /// /// Set to when a concurrent close is already in flight; the caller /// should treat the session as already being closed (same semantics as /// ). /// /// /// when the state was flipped to Closing and the caller /// should proceed with teardown; when the session is already /// closed OR is no longer eligible (a subscriber re-attached between the eligibility /// check in the sweep loop and this call — the reconnect won the race and the session /// should be left open). /// /// /// /// Race: CloseExpiredLeasesAsync evaluates / /// outside the close lock, then calls /// which takes _closeLock. A client can call /// in between, clearing _detachedAtUtc and /// incrementing _activeEventSubscriberCount — the session is no longer expired. /// This method re-checks eligibility atomically under _syncRoot before /// committing to Closing, so a reattach that wins the race leaves the session /// in Ready and usable. /// /// internal bool TryBeginCloseIfExpired(DateTimeOffset now, out bool alreadyClosing) { lock (_syncRoot) { if (_state is SessionState.Closed) { alreadyClosing = _closeStarted; return false; } // Re-verify eligibility atomically. If a subscriber reattached between the sweep's // eligibility check and this point, neither condition holds and we decline. bool eligible = IsLeaseExpiredCore(now) || IsFaultedReapableCore(now) || IsDetachGraceExpiredCore(now); if (!eligible) { alreadyClosing = false; return false; } alreadyClosing = _closeStarted; _closeStarted = true; _state = SessionState.Closing; return true; } } // Lock-free (must be called under _syncRoot) helpers used by TryBeginCloseIfExpired. private bool IsLeaseExpiredCore(DateTimeOffset now) => _activeEventSubscriberCount == 0 && _leaseExpiresAt is not null && _leaseExpiresAt <= now; private bool IsDetachGraceExpiredCore(DateTimeOffset now) => _detachGrace > TimeSpan.Zero && _activeEventSubscriberCount == 0 && _detachedAtUtc is not null && now - _detachedAtUtc.Value >= _detachGrace; private bool IsFaultedReapableCore(DateTimeOffset now) => _state is SessionState.Faulted && (_faultedGrace <= TimeSpan.Zero || _faultedAtUtc is null || now - _faultedAtUtc.Value >= _faultedGrace); // Final terminal transition; under _syncRoot to keep _state writes single-lock. // Closed is unconditionally terminal — TransitionTo refuses to overwrite it — // so we don't need to re-check the precondition here. private void MarkClosed() { lock (_syncRoot) { _state = SessionState.Closed; } } /// /// Terminates the worker process immediately. /// /// Reason for killing the worker. public void KillWorker(string reason) { _workerClient?.Kill(reason); TransitionTo(SessionState.Closed); } /// /// Terminates the worker process immediately while holding the per-session /// close lock so concurrent close/kill callers serialize. Returns the /// session state observed at the start of the call so the caller can /// dedup metric accounting (e.g. only record SessionClosed when /// the session was not already closed). /// /// /// Mirrors 's use of _closeLock so that /// a Close in flight from one caller and a Kill from another do not /// race on the "was the session already closed" observation that /// drives metric increments. /// /// Reason for killing the worker. /// Cancellation token. /// true if the session was already when the lock was acquired; otherwise false. public async ValueTask KillWorkerWithCloseGateAsync( string reason, CancellationToken cancellationToken) { await _closeLock.WaitAsync(cancellationToken).ConfigureAwait(false); try { bool wasClosed; lock (_syncRoot) { wasClosed = _state == SessionState.Closed; } _workerClient?.Kill(reason); TransitionTo(SessionState.Closed); return wasClosed; } finally { _closeLock.Release(); } } /// /// Disposes the session and frees associated resources. /// /// /// Acquires _closeLock once before disposing so an in-flight /// finishes before the semaphore is released and /// reclaimed. Without this gate, the in-flight close's _closeLock.Release() /// would race the dispose and raise . /// The acquire is best-effort: a non-cancellable wait that swallows /// so double-dispose still completes. /// /// A task that represents the asynchronous operation. public async ValueTask DisposeAsync() { try { // CancellationToken.None — disposal must not be cancelled, and a misbehaving // close path that never releases would have to be torn down by the worker // shutdown timeout long before we reach here. await _closeLock.WaitAsync(CancellationToken.None).ConfigureAwait(false); try { // Hand the slot back so the semaphore's internal counter is consistent // for any contemporaneous waiter, then dispose. Once disposed, every // subsequent WaitAsync / Release will throw — but DisposeAsync's contract // is "no concurrent close after this point", which SessionManager honors. _closeLock.Release(); } catch (ObjectDisposedException) { } } catch (ObjectDisposedException) { // Already disposed (e.g. double-dispose); nothing to gate on. } try { _closeLock.Dispose(); } catch (ObjectDisposedException) { } // Stop the internal dashboard mirror first: cancel its loop, dispose its lease (which // unregisters its internal distributor subscriber and completes its channel), and // await the loop task. Done BEFORE disposing the distributor and worker client — like // the distributor itself — so the mirror is no longer reading the pump when the pump // and its source (the worker client) tear down. IEventSubscriberLease? dashboardLease; Task? dashboardTask; CancellationTokenSource? dashboardCts; lock (_syncRoot) { dashboardLease = _dashboardMirrorLease; dashboardTask = _dashboardMirrorTask; dashboardCts = _dashboardMirrorCts; _dashboardMirrorLease = null; _dashboardMirrorTask = null; _dashboardMirrorCts = null; } if (dashboardCts is not null) { await dashboardCts.CancelAsync().ConfigureAwait(false); } dashboardLease?.Dispose(); if (dashboardTask is not null) { try { await dashboardTask.ConfigureAwait(false); } catch (Exception) { // The mirror loop swallows its own faults; any escape here must not block // disposal. The loop has stopped, which is all teardown requires. } } dashboardCts?.Dispose(); // Stop the event pump and complete every subscriber channel before tearing down the // worker client (the pump's source). DisposeAsync is the single session teardown // point (SessionManager.RemoveSessionAsync awaits it after close), so awaiting it // here guarantees the distributor's pump task is observed and subscribers are // completed rather than left dangling. SessionEventDistributor? distributor; lock (_syncRoot) { distributor = _eventDistributor; _eventDistributor = null; } if (distributor is not null) { await distributor.DisposeAsync().ConfigureAwait(false); } if (_workerClient is not null) { await _workerClient.DisposeAsync().ConfigureAwait(false); } } private async Task> InvokeBulkAsync( MxCommand command, Func payloadAccessor, CancellationToken cancellationToken) { MxCommandReply reply = await InvokeBulkInternalAsync(command, cancellationToken).ConfigureAwait(false); return payloadAccessor(reply)?.Results.ToArray() ?? []; } private async Task> InvokeBulkWriteAsync( MxCommand command, Func payloadAccessor, CancellationToken cancellationToken) { MxCommandReply reply = await InvokeBulkInternalAsync(command, cancellationToken).ConfigureAwait(false); return payloadAccessor(reply)?.Results.ToArray() ?? []; } private async Task> InvokeBulkReadAsync( MxCommand command, Func payloadAccessor, CancellationToken cancellationToken) { MxCommandReply reply = await InvokeBulkInternalAsync(command, cancellationToken).ConfigureAwait(false); return payloadAccessor(reply)?.Results.ToArray() ?? []; } // Single round-trip + protocol-status check shared by every bulk variant. // Callers project the typed reply payload out via their own accessor — the // outer envelope handling is identical across SubscribeResult-based bulks, // BulkWriteResult-based writes, and BulkReadResult-based reads. private async Task InvokeBulkInternalAsync( MxCommand command, CancellationToken cancellationToken) { WorkerCommandReply workerReply = await InvokeAsync( new WorkerCommand { Command = command }, cancellationToken) .ConfigureAwait(false); MxCommandReply reply = workerReply.Reply ?? new MxCommandReply { ProtocolStatus = new ProtocolStatus { Code = ProtocolStatusCode.ProtocolViolation, Message = "Worker command reply did not contain a public reply payload.", }, }; if (reply.ProtocolStatus?.Code is not ProtocolStatusCode.Ok) { string message = reply.ProtocolStatus?.Message ?? reply.DiagnosticMessage; throw new SessionManagerException( SessionManagerErrorCode.SessionNotReady, string.IsNullOrWhiteSpace(message) ? "Bulk MXAccess command failed." : message); } return reply; } /// /// Bounded, opt-in async variant of the fail-fast readiness check. When the /// session is but the worker has transiently diverged /// to a non-terminal state (/ /// ) and the configured worker-ready wait timeout /// is positive, this polls (outside _syncRoot) until the worker reaches /// or the deadline elapses, re-evaluating the /// fast-path/fail-fast decision under the lock on each poll. Terminal worker states, a /// missing worker, or a non-Ready session fail fast immediately. With the default /// timeout of zero this behaves byte-for-byte like the synchronous fail-fast path: no /// await, no delay. /// /// Token to cancel the wait. /// The worker client once both the session and worker are Ready. private async Task GetReadyWorkerClientAsync(CancellationToken cancellationToken) { const int pollIntervalMs = 25; string? failureMessage; lock (_syncRoot) { IWorkerClient? ready = EvaluateReadyUnderLock(out failureMessage); if (ready is not null) { return ready; } // Only transient (non-terminal) worker states with a positive wait timeout fall // through to the bounded wait loop. Everything else (terminal worker, no worker, // session not Ready, or a zero timeout) fails fast right here under the lock. When // the worker is merely transient (failureMessage is null) but the wait is disabled, // build the both-states diagnostic so the zero-timeout path is byte-for-byte the // original fail-fast message. if (failureMessage is not null || _workerReadyWaitTimeout <= TimeSpan.Zero) { throw new SessionManagerException( SessionManagerErrorCode.SessionNotReady, failureMessage ?? BuildNotReadyMessage()); } } DateTimeOffset deadline = _eventStreaming.TimeProvider.GetUtcNow() + _workerReadyWaitTimeout; while (true) { await Task.Delay( TimeSpan.FromMilliseconds(pollIntervalMs), _eventStreaming.TimeProvider, cancellationToken) .ConfigureAwait(false); lock (_syncRoot) { IWorkerClient? ready = EvaluateReadyUnderLock(out failureMessage); if (ready is not null) { return ready; } // A terminal worker / missing worker / non-Ready session surfaced while we // waited: fail fast immediately rather than burning the rest of the deadline. if (failureMessage is not null) { throw new SessionManagerException(SessionManagerErrorCode.SessionNotReady, failureMessage); } } if (_eventStreaming.TimeProvider.GetUtcNow() >= deadline) { lock (_syncRoot) { IWorkerClient? ready = EvaluateReadyUnderLock(out failureMessage); if (ready is not null) { return ready; } throw new SessionManagerException( SessionManagerErrorCode.SessionNotReady, failureMessage ?? BuildNotReadyMessage()); } } } } /// /// Evaluates readiness while the caller already holds _syncRoot. Returns the /// worker client when both the session and worker are /// (with set to ). Returns /// together with the both-states diagnostic in /// when the worker is in a terminal state /// (// /// ), there is no worker, or the session is not /// . Returns with a /// when the session is /// Ready but the worker is in a transient state /// (/) — /// the signal for the async path to keep waiting. /// /// /// The fail-fast both-states diagnostic when readiness cannot succeed, or /// for the keep-waiting (transient) signal. /// /// The ready worker client, or . private IWorkerClient? EvaluateReadyUnderLock(out string? failureMessage) { if (_state == SessionState.Ready && _workerClient?.State == WorkerClientState.Ready) { failureMessage = null; return _workerClient; } // Keep-waiting signal: session is Ready and the worker is merely transient. if (_state == SessionState.Ready && _workerClient is { State: WorkerClientState.Handshaking or WorkerClientState.Created }) { failureMessage = null; return null; } failureMessage = BuildNotReadyMessage(); return null; } /// Builds the both-states not-ready diagnostic (must be called under _syncRoot). /// The diagnostic message surfacing both the session and worker states. private string BuildNotReadyMessage() { string workerState = _workerClient is null ? "" : _workerClient.State.ToString(); return $"Session {SessionId} is not ready. Session state is {_state}; worker state is {workerState}."; } private void TrackItem( int serverHandle, int itemHandle, string tagAddress) { if (itemHandle == 0 || string.IsNullOrWhiteSpace(tagAddress)) { return; } _items[(serverHandle, itemHandle)] = new SessionItemRegistration(serverHandle, itemHandle, tagAddress); } private void TrackBulkItems(BulkSubscribeReply reply) { foreach (SubscribeResult result in reply.Results) { if (result.WasSuccessful) { TrackItem(result.ServerHandle, result.ItemHandle, result.TagAddress); } } } private void RemoveItems( int serverHandle, IEnumerable itemHandles) { foreach (int itemHandle in itemHandles) { _items.Remove((serverHandle, itemHandle)); } } private void DetachEventSubscriber() { lock (_syncRoot) { // Assert in debug so a genuine double-decrement (a logic error) surfaces // loudly; the clamp below keeps release builds safe if it somehow fires. Debug.Assert(_activeEventSubscriberCount > 0, "DetachEventSubscriber called with _activeEventSubscriberCount already at 0 — possible double-dispose."); if (_activeEventSubscriberCount > 0) { _activeEventSubscriberCount--; } // When the LAST external subscriber drops and detach-grace is enabled, retain the // session instead of letting it linger only on the (long) lease: stamp the detached // time so the lease monitor can close it once the grace window elapses. The session // stays in its current (Ready) state and remains usable, so a reconnecting subscriber // re-attaches normally. The gateway-owned internal dashboard subscriber is // NOT counted in _activeEventSubscriberCount (it registers on the distributor with // isInternal: true), so a session whose only remaining subscriber is the dashboard // mirror still enters grace. Only stamp while the session is alive — once // Closing/Closed/Faulted there is nothing to retain. This is the detach→grace-start // transition; it shares _syncRoot with the reattach→grace-cancel write above and the // sweeper's IsDetachGraceExpired read, so the three serialize. // Only stamp a detach that mirrors a prior SUCCESSFUL attach. The attach catch path // calls this same method to roll back a reserved slot when the FIRST attach failed // before any subscriber registered; that never-subscribed session must not enter the // grace window. if (_everHadEventSubscriber && _detachGrace > TimeSpan.Zero && _activeEventSubscriberCount == 0 && _state is not (SessionState.Closing or SessionState.Closed or SessionState.Faulted)) { _detachedAtUtc = _eventStreaming.TimeProvider.GetUtcNow(); } } } private sealed class EventSubscriberLease(GatewaySession session, IEventSubscriberLease distributorLease) : IEventSubscriberLease { // 0 = live, 1 = disposed. Interlocked so concurrent stream-completion + // client-cancellation paths cannot both call DetachEventSubscriber and // double-decrement _activeEventSubscriberCount to -1. private int _leaseDisposed; /// public System.Threading.Channels.ChannelReader Reader => distributorLease.Reader; /// /// Disposes the lease: unregisters this subscriber from the distributor (completing /// its channel) and decrements the session's active-subscriber count. Ordering is /// not significant — the count guard and the distributor registration are /// independent — but both must run exactly once. /// public void Dispose() { if (Interlocked.Exchange(ref _leaseDisposed, 1) == 0) { distributorLease.Dispose(); session.DetachEventSubscriber(); } } } }