7bdb6d00ec
Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
130 lines
5.3 KiB
C#
130 lines
5.3 KiB
C#
using System.Collections.Concurrent;
|
|
using System.Threading.Channels;
|
|
|
|
namespace ZB.MOM.WW.OtOpcUa.Runtime.Telemetry;
|
|
|
|
/// <summary>
|
|
/// Default <see cref="ITelemetryLocalHub"/>: a process-wide singleton on driver-role nodes.
|
|
/// </summary>
|
|
/// <remarks>
|
|
/// <para>
|
|
/// <b>No DPS.</b> This hub is fed only by this node's own in-process producers (a later Phase-5
|
|
/// task taps them). It never subscribes to cluster DistributedPubSub — see the invariant on
|
|
/// <see cref="ITelemetryLocalHub"/>.
|
|
/// </para>
|
|
/// <para>
|
|
/// <b>Snapshot-then-attach ordering (no lost delta).</b> <see cref="Emit"/> updates the snapshot
|
|
/// cache and fans out under a single gate; <see cref="Subscribe"/> writes the cached snapshots
|
|
/// into the new channel and only THEN adds it to the subscriber set — all under that same gate.
|
|
/// Because both operations are serialized by the gate, an <see cref="Emit"/> racing a
|
|
/// <see cref="Subscribe"/> is resolved one of two ways, both correct: it runs entirely before the
|
|
/// subscribe (its value is in the snapshot, and it is not delivered live because the channel is
|
|
/// not yet attached), or entirely after (it is delivered live, strictly after the snapshot). No
|
|
/// item is dropped or duplicated across the boundary. Fan-out is a non-blocking
|
|
/// <see cref="ChannelWriter{T}.TryWrite"/> into bounded/DropOldest channels, so holding the gate
|
|
/// never blocks on a slow consumer.
|
|
/// </para>
|
|
/// </remarks>
|
|
public sealed class TelemetryLocalHub : ITelemetryLocalHub
|
|
{
|
|
private readonly object _gate = new();
|
|
|
|
// Plain Dictionary, NOT ConcurrentDictionary: every access is already under _gate, and
|
|
// ConcurrentDictionary.Values would materialize a fresh List (and take its own internal lock) on
|
|
// every Emit — the hot node-actor-thread path. Under the gate a plain Dictionary is safe and
|
|
// allocation-free to iterate.
|
|
private readonly Dictionary<Guid, Channel<TelemetryItem>> _subscribers = new();
|
|
|
|
// Snapshot last-value caches. Only ever mutated under _gate (kept ConcurrentDictionary so a
|
|
// Subscribe reading them under the gate is trivially safe even against any future lock-free read).
|
|
// TODO(mesh-phase5+): evict snapshot cache entries on driver-instance removal (needs a
|
|
// driver-lifecycle hook); pre-production, accepted for now. Until then a decommissioned driver
|
|
// instance's last-known Health/Resilience state replays to new subscribers forever.
|
|
private readonly ConcurrentDictionary<string, TelemetryItem.Health> _healthCache = new();
|
|
private readonly ConcurrentDictionary<(string InstanceId, string HostName), TelemetryItem.Resilience> _resilienceCache = new();
|
|
|
|
/// <inheritdoc />
|
|
public void Emit(TelemetryItem item)
|
|
{
|
|
ArgumentNullException.ThrowIfNull(item);
|
|
|
|
lock (_gate)
|
|
{
|
|
switch (item)
|
|
{
|
|
case TelemetryItem.Health h:
|
|
_healthCache[h.E.DriverInstanceId] = h;
|
|
break;
|
|
case TelemetryItem.Resilience r:
|
|
_resilienceCache[(r.E.DriverInstanceId, r.E.HostName)] = r;
|
|
break;
|
|
// Alarm / Script are append-style — never cached.
|
|
}
|
|
|
|
foreach (var channel in _subscribers.Values)
|
|
channel.Writer.TryWrite(item); // bounded + DropOldest ⇒ non-blocking, sheds oldest when full
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public ITelemetrySubscription Subscribe(int boundedCapacity)
|
|
{
|
|
if (boundedCapacity < 1)
|
|
throw new ArgumentOutOfRangeException(nameof(boundedCapacity), boundedCapacity, "Capacity must be positive.");
|
|
|
|
var channel = Channel.CreateBounded<TelemetryItem>(new BoundedChannelOptions(boundedCapacity)
|
|
{
|
|
FullMode = BoundedChannelFullMode.DropOldest,
|
|
SingleReader = true,
|
|
SingleWriter = false,
|
|
});
|
|
|
|
var id = Guid.NewGuid();
|
|
lock (_gate)
|
|
{
|
|
// Snapshot FIRST (Health then Resilience), attach SECOND — both under the gate so a
|
|
// concurrent Emit cannot slip a delta between the snapshot read and the attach.
|
|
foreach (var health in _healthCache.Values)
|
|
channel.Writer.TryWrite(health);
|
|
foreach (var resilience in _resilienceCache.Values)
|
|
channel.Writer.TryWrite(resilience);
|
|
|
|
_subscribers[id] = channel;
|
|
}
|
|
|
|
return new Subscription(this, id, channel);
|
|
}
|
|
|
|
private void Detach(Guid id)
|
|
{
|
|
lock (_gate)
|
|
{
|
|
if (_subscribers.Remove(id, out var channel))
|
|
channel.Writer.TryComplete();
|
|
}
|
|
}
|
|
|
|
private sealed class Subscription : ITelemetrySubscription
|
|
{
|
|
private readonly TelemetryLocalHub _hub;
|
|
private readonly Guid _id;
|
|
private int _disposed;
|
|
|
|
public Subscription(TelemetryLocalHub hub, Guid id, Channel<TelemetryItem> channel)
|
|
{
|
|
_hub = hub;
|
|
_id = id;
|
|
Reader = channel.Reader;
|
|
}
|
|
|
|
public ChannelReader<TelemetryItem> Reader { get; }
|
|
|
|
public void Dispose()
|
|
{
|
|
if (Interlocked.Exchange(ref _disposed, 1) != 0)
|
|
return; // idempotent
|
|
_hub.Detach(_id);
|
|
}
|
|
}
|
|
}
|