dc9c0c950c
Apply the ZB.MOM.WW. prefix to all gateway-side projects, folders,
.csproj/.sln contents, C# namespaces, using directives, generated proto
C# (csharp_namespace + checked-in generated files), InternalsVisibleTo
attributes, project-name string literals (LoadProject, .sln lookups,
worker exe paths, staticwebassets manifest), and the install/script/doc
references that point at any of the above. Migrate the solution from
.sln to .slnx via `dotnet sln migrate` and delete the old file.
External-runtime identifiers are intentionally NOT prefixed so external
configuration keeps working:
- GatewayMetrics.cs MeterName ("MxGateway.Server")
- DashboardAuthenticationDefaults Scheme/Policy ("MxGateway.Dashboard")
- GatewayRequestLoggingMiddleware logger category ("MxGateway.Request")
- StaRuntime thread name ("MxGateway.Worker.STA")
- appsettings.json root section "MxGateway" + env-var prefix
MxGateway__... and secret-name MxGateway:ApiKeyPepper
- C:\ProgramData\MxGateway\ data dir paths
Also fixes two tests that were not rename-related but became visible
while validating the rename:
- WorkerLiveMxAccessSmokeTests.ShutDownAsync: cancellation that the
gateway service correctly maps to RpcException(Cancelled) per gRPC
convention was being misclassified as a stream fault. Added a sibling
catch on RpcException with StatusCode.Cancelled.
- IntegrationTestEnvironment.ResolveRepositoryRoot: extracted IsRepositoryRoot
and made it accept either a .git marker OR a .sln/.slnx next to src/
so the worker-exe walker works in non-git working copies.
clients/proto/proto-inputs.json's protoRoot updated to point at
src/ZB.MOM.WW.MxGateway.Contracts/Protos.
Verified by `dotnet build` and a full `dotnet test` of the .slnx with
MXGATEWAY_RUN_LIVE_{MXACCESS,LDAP,GALAXY}_TESTS=1:
Tests: 472/472 pass
Worker.Tests: 280/280 pass (4 dev-rig [Fact(Skip=...)] skipped)
IntegrationTests: 18/18 pass
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
83 lines
2.9 KiB
C#
83 lines
2.9 KiB
C#
using System.Collections.Concurrent;
|
|
using System.Runtime.CompilerServices;
|
|
using System.Threading.Channels;
|
|
|
|
namespace ZB.MOM.WW.MxGateway.Server.Galaxy;
|
|
|
|
/// <summary>
|
|
/// Channel-based fan-out of Galaxy deploy events to streaming gRPC subscribers. Each
|
|
/// subscriber gets a private bounded channel so a slow client cannot back-pressure
|
|
/// other subscribers or the publisher. When a subscriber's channel is full the oldest
|
|
/// event is dropped — clients use the sequence field to detect gaps.
|
|
/// </summary>
|
|
/// <summary>
|
|
/// Publishes Galaxy deploy events to streaming gRPC subscribers via private bounded channels.
|
|
/// </summary>
|
|
public sealed class GalaxyDeployNotifier : IGalaxyDeployNotifier
|
|
{
|
|
private const int SubscriberQueueCapacity = 16;
|
|
|
|
private readonly ConcurrentDictionary<Guid, Channel<GalaxyDeployEventInfo>> _subscribers = new();
|
|
private GalaxyDeployEventInfo? _latest;
|
|
|
|
/// <summary>
|
|
/// The most recent deploy event, or null if none has been published.
|
|
/// </summary>
|
|
public GalaxyDeployEventInfo? Latest => Volatile.Read(ref _latest);
|
|
|
|
/// <inheritdoc />
|
|
public void Publish(GalaxyDeployEventInfo info)
|
|
{
|
|
ArgumentNullException.ThrowIfNull(info);
|
|
|
|
Volatile.Write(ref _latest, info);
|
|
|
|
foreach (Channel<GalaxyDeployEventInfo> channel in _subscribers.Values)
|
|
{
|
|
// BoundedChannelFullMode.DropOldest -> writes never wait; we only fail if the
|
|
// channel was completed by the subscriber side, which we ignore.
|
|
channel.Writer.TryWrite(info);
|
|
}
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async IAsyncEnumerable<GalaxyDeployEventInfo> SubscribeAsync(
|
|
[EnumeratorCancellation] CancellationToken cancellationToken)
|
|
{
|
|
Guid subscriberId = Guid.NewGuid();
|
|
Channel<GalaxyDeployEventInfo> channel = Channel.CreateBounded<GalaxyDeployEventInfo>(
|
|
new BoundedChannelOptions(SubscriberQueueCapacity)
|
|
{
|
|
FullMode = BoundedChannelFullMode.DropOldest,
|
|
SingleReader = true,
|
|
SingleWriter = false,
|
|
});
|
|
|
|
_subscribers[subscriberId] = channel;
|
|
|
|
// Bootstrap: emit the latest known event so subscribers don't need to wait for
|
|
// the next deploy to know current state.
|
|
GalaxyDeployEventInfo? bootstrap = Volatile.Read(ref _latest);
|
|
if (bootstrap is not null)
|
|
{
|
|
channel.Writer.TryWrite(bootstrap);
|
|
}
|
|
|
|
try
|
|
{
|
|
while (await channel.Reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false))
|
|
{
|
|
while (channel.Reader.TryRead(out GalaxyDeployEventInfo? next))
|
|
{
|
|
yield return next;
|
|
}
|
|
}
|
|
}
|
|
finally
|
|
{
|
|
_subscribers.TryRemove(subscriberId, out _);
|
|
channel.Writer.TryComplete();
|
|
}
|
|
}
|
|
}
|