feat(mesh-phase5): dark-switch central telemetry ingest (Dps bridges | Grpc dialer)
Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
This commit is contained in:
@@ -1,7 +1,14 @@
|
||||
using Akka.Actor;
|
||||
using Akka.Hosting;
|
||||
using Microsoft.AspNetCore.SignalR;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
using ZB.MOM.WW.OtOpcUa.AdminUI.Telemetry;
|
||||
using ZB.MOM.WW.OtOpcUa.Cluster;
|
||||
using ZB.MOM.WW.OtOpcUa.Configuration;
|
||||
using ZB.MOM.WW.OtOpcUa.ControlPlane.Telemetry;
|
||||
|
||||
namespace ZB.MOM.WW.OtOpcUa.AdminUI.Hubs;
|
||||
|
||||
@@ -12,6 +19,7 @@ public static class HubServiceCollectionExtensions
|
||||
public const string ScriptLogSignalRBridgeName = "script-log-signalr-bridge";
|
||||
public const string DriverStatusSignalRBridgeName = "driver-status-signalr-bridge";
|
||||
public const string DriverResilienceStatusBridgeName = "driver-resilience-status-bridge";
|
||||
public const string TelemetryDialSupervisorName = "telemetry-dial-supervisor";
|
||||
|
||||
/// <summary>
|
||||
/// Registers the in-process live-push services the AdminUI's Blazor Server panels read
|
||||
@@ -54,30 +62,69 @@ public static class HubServiceCollectionExtensions
|
||||
{
|
||||
builder.WithActors((system, registry, resolver) =>
|
||||
{
|
||||
// Fleet-status always stays on DPS in both modes (deferred / out of Phase 5 scope) — it is
|
||||
// never gated by the telemetry dark switch.
|
||||
var fleetHub = resolver.GetService<IHubContext<FleetStatusHub>>();
|
||||
var fleetBridge = system.ActorOf(FleetStatusSignalRBridge.Props(fleetHub), FleetStatusSignalRBridgeName);
|
||||
registry.Register<FleetStatusSignalRBridgeKey>(fleetBridge);
|
||||
|
||||
var alertHub = resolver.GetService<IHubContext<AlertHub>>();
|
||||
// The four telemetry sinks are registered identically in both modes; only the upstream that
|
||||
// feeds them swaps (DPS bridges vs. the gRPC dial supervisor). Resolve them once, above the
|
||||
// branch, so both paths feed the SAME singletons the Blazor panels read.
|
||||
var alertBroadcaster = resolver.GetService<IInProcessBroadcaster<Commons.Messages.Alerts.AlarmTransitionEvent>>();
|
||||
var alertBridge = system.ActorOf(AlertSignalRBridge.Props(alertHub, alertBroadcaster), AlertSignalRBridgeName);
|
||||
registry.Register<AlertSignalRBridgeKey>(alertBridge);
|
||||
|
||||
var scriptLogHub = resolver.GetService<IHubContext<ScriptLogHub>>();
|
||||
var scriptLogBroadcaster = resolver.GetService<IInProcessBroadcaster<Commons.Messages.Logging.ScriptLogEntry>>();
|
||||
var scriptLogBridge = system.ActorOf(ScriptLogSignalRBridge.Props(scriptLogHub, scriptLogBroadcaster), ScriptLogSignalRBridgeName);
|
||||
registry.Register<ScriptLogSignalRBridgeKey>(scriptLogBridge);
|
||||
|
||||
var driverStatusHub = resolver.GetService<IHubContext<DriverStatusHub>>();
|
||||
var driverStatusStore = resolver.GetService<IDriverStatusSnapshotStore>();
|
||||
var driverStatusBridge = system.ActorOf(DriverStatusSignalRBridge.Props(driverStatusHub, driverStatusStore), DriverStatusSignalRBridgeName);
|
||||
registry.Register<DriverStatusSignalRBridgeKey>(driverStatusBridge);
|
||||
|
||||
// Resilience-status bridge: DPS topic -> in-process store (no SignalR hub — the panel reads
|
||||
// the store directly, and resilience has no browser-JS consumer).
|
||||
var resilienceStore = resolver.GetService<IDriverResilienceStatusStore>();
|
||||
var resilienceBridge = system.ActorOf(DriverResilienceStatusBridge.Props(resilienceStore), DriverResilienceStatusBridgeName);
|
||||
registry.Register<DriverResilienceStatusBridgeKey>(resilienceBridge);
|
||||
|
||||
// Phase 5 dark switch. Absent options ⇒ Dps (today's behaviour); case-insensitive.
|
||||
var telemetryMode = resolver.GetService<IOptions<TelemetryDialOptions>>()?.Value.Mode
|
||||
?? TelemetryDialOptions.ModeDps;
|
||||
|
||||
if (string.Equals(telemetryMode, TelemetryDialOptions.ModeGrpc, StringComparison.OrdinalIgnoreCase))
|
||||
{
|
||||
// Grpc: central dials each enabled node's dedicated telemetry stream, feeding the SAME
|
||||
// four sinks the DPS bridges feed. No DPS telemetry bridges are spawned.
|
||||
var options = resolver.GetService<IOptions<TelemetryDialOptions>>()!.Value;
|
||||
var dbFactory = resolver.GetService<IDbContextFactory<OtOpcUaConfigDbContext>>();
|
||||
var loggerFactory = resolver.GetService<ILoggerFactory>();
|
||||
|
||||
var nodeSource = TelemetryNodeSource.Create(
|
||||
dbFactory!, loggerFactory!.CreateLogger(typeof(TelemetryNodeSource).FullName!));
|
||||
var dialLoop = TelemetryNodeSource.CreateDialLoop(
|
||||
options.ApiKey, loggerFactory.CreateLogger<TelemetryStreamClient>());
|
||||
|
||||
var supervisor = system.ActorOf(
|
||||
TelemetryDialSupervisor.Props(
|
||||
nodeSource,
|
||||
dialLoop,
|
||||
alertBroadcaster,
|
||||
scriptLogBroadcaster,
|
||||
driverStatusStore,
|
||||
resilienceStore,
|
||||
options),
|
||||
TelemetryDialSupervisorName);
|
||||
registry.Register<TelemetryDialSupervisorKey>(supervisor);
|
||||
}
|
||||
else
|
||||
{
|
||||
// Dps (default): the four DPS bridges subscribe their mesh-wide topics and feed the sinks.
|
||||
var alertHub = resolver.GetService<IHubContext<AlertHub>>();
|
||||
var alertBridge = system.ActorOf(AlertSignalRBridge.Props(alertHub, alertBroadcaster), AlertSignalRBridgeName);
|
||||
registry.Register<AlertSignalRBridgeKey>(alertBridge);
|
||||
|
||||
var scriptLogHub = resolver.GetService<IHubContext<ScriptLogHub>>();
|
||||
var scriptLogBridge = system.ActorOf(ScriptLogSignalRBridge.Props(scriptLogHub, scriptLogBroadcaster), ScriptLogSignalRBridgeName);
|
||||
registry.Register<ScriptLogSignalRBridgeKey>(scriptLogBridge);
|
||||
|
||||
var driverStatusHub = resolver.GetService<IHubContext<DriverStatusHub>>();
|
||||
var driverStatusBridge = system.ActorOf(DriverStatusSignalRBridge.Props(driverStatusHub, driverStatusStore), DriverStatusSignalRBridgeName);
|
||||
registry.Register<DriverStatusSignalRBridgeKey>(driverStatusBridge);
|
||||
|
||||
// Resilience-status bridge: DPS topic -> in-process store (no SignalR hub — the panel reads
|
||||
// the store directly, and resilience has no browser-JS consumer).
|
||||
var resilienceBridge = system.ActorOf(DriverResilienceStatusBridge.Props(resilienceStore), DriverResilienceStatusBridgeName);
|
||||
registry.Register<DriverResilienceStatusBridgeKey>(resilienceBridge);
|
||||
}
|
||||
});
|
||||
return builder;
|
||||
}
|
||||
@@ -89,3 +136,6 @@ public sealed class AlertSignalRBridgeKey { }
|
||||
public sealed class ScriptLogSignalRBridgeKey { }
|
||||
public sealed class DriverStatusSignalRBridgeKey { }
|
||||
public sealed class DriverResilienceStatusBridgeKey { }
|
||||
|
||||
/// <summary>Marker key for <see cref="ActorRegistry"/> lookup of the Grpc-mode telemetry dial supervisor.</summary>
|
||||
public sealed class TelemetryDialSupervisorKey { }
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
using Akka.Actor;
|
||||
using Akka.Hosting;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Options;
|
||||
using Shouldly;
|
||||
using Xunit;
|
||||
using ZB.MOM.WW.OtOpcUa.AdminUI.Hubs;
|
||||
using ZB.MOM.WW.OtOpcUa.Cluster;
|
||||
using ZB.MOM.WW.OtOpcUa.Configuration;
|
||||
|
||||
namespace ZB.MOM.WW.OtOpcUa.AdminUI.Tests.Hubs;
|
||||
|
||||
/// <summary>
|
||||
/// Verifies the Phase 5 dark switch in <see cref="HubServiceCollectionExtensions.WithOtOpcUaSignalRBridges"/>:
|
||||
/// <c>TelemetryDial:Mode = Dps</c> (default) spawns the four telemetry DPS bridge actors and NOT the
|
||||
/// gRPC dial supervisor; <c>Grpc</c> spawns the dial supervisor and NONE of the four telemetry
|
||||
/// bridges. The fleet-status bridge (deferred / out of Phase 5 scope) is spawned in BOTH modes.
|
||||
/// The gate is proven on a real <see cref="ActorRegistry"/> from a started Akka host.
|
||||
/// </summary>
|
||||
public sealed class TelemetryModeWiringTests
|
||||
{
|
||||
/// <summary>Dps mode: the four telemetry bridges + fleet bridge are registered; the dial supervisor is not.</summary>
|
||||
[Fact]
|
||||
public async Task Dps_mode_spawns_the_four_telemetry_bridges_and_not_the_dial_supervisor()
|
||||
{
|
||||
using var host = BuildBridgeHost(TelemetryDialOptions.ModeDps);
|
||||
await host.StartAsync();
|
||||
try
|
||||
{
|
||||
var registry = host.Services.GetRequiredService<ActorRegistry>();
|
||||
|
||||
registry.TryGet<FleetStatusSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<AlertSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<ScriptLogSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<DriverStatusSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<DriverResilienceStatusBridgeKey>(out _).ShouldBeTrue();
|
||||
|
||||
registry.TryGet<TelemetryDialSupervisorKey>(out _).ShouldBeFalse();
|
||||
}
|
||||
finally
|
||||
{
|
||||
await host.StopAsync();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>Grpc mode: the dial supervisor + fleet bridge are registered; none of the four telemetry bridges are.</summary>
|
||||
[Fact]
|
||||
public async Task Grpc_mode_spawns_the_dial_supervisor_and_none_of_the_four_telemetry_bridges()
|
||||
{
|
||||
using var host = BuildBridgeHost(TelemetryDialOptions.ModeGrpc);
|
||||
await host.StartAsync();
|
||||
try
|
||||
{
|
||||
var registry = host.Services.GetRequiredService<ActorRegistry>();
|
||||
|
||||
registry.TryGet<FleetStatusSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<TelemetryDialSupervisorKey>(out _).ShouldBeTrue();
|
||||
|
||||
registry.TryGet<AlertSignalRBridgeKey>(out _).ShouldBeFalse();
|
||||
registry.TryGet<ScriptLogSignalRBridgeKey>(out _).ShouldBeFalse();
|
||||
registry.TryGet<DriverStatusSignalRBridgeKey>(out _).ShouldBeFalse();
|
||||
registry.TryGet<DriverResilienceStatusBridgeKey>(out _).ShouldBeFalse();
|
||||
}
|
||||
finally
|
||||
{
|
||||
await host.StopAsync();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>Unset (default) mode is treated as Dps — the four bridges spawn, the supervisor does not.</summary>
|
||||
[Fact]
|
||||
public async Task Default_mode_is_dps()
|
||||
{
|
||||
using var host = BuildBridgeHost(mode: null);
|
||||
await host.StartAsync();
|
||||
try
|
||||
{
|
||||
var registry = host.Services.GetRequiredService<ActorRegistry>();
|
||||
registry.TryGet<AlertSignalRBridgeKey>(out _).ShouldBeTrue();
|
||||
registry.TryGet<TelemetryDialSupervisorKey>(out _).ShouldBeFalse();
|
||||
}
|
||||
finally
|
||||
{
|
||||
await host.StopAsync();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>Builds an admin-role host that runs <c>WithOtOpcUaSignalRBridges</c> under the given telemetry mode.</summary>
|
||||
/// <param name="mode">The <c>TelemetryDial:Mode</c> value, or <see langword="null"/> to leave it at its default.</param>
|
||||
private static IHost BuildBridgeHost(string? mode)
|
||||
=> Host.CreateDefaultBuilder()
|
||||
.ConfigureServices((_, services) =>
|
||||
{
|
||||
services.AddSignalR();
|
||||
services.AddOtOpcUaDriverStatusServices();
|
||||
services.AddSingleton<IDbContextFactory<OtOpcUaConfigDbContext>>(
|
||||
new InMemoryConfigDbFactory(Guid.NewGuid().ToString("N")));
|
||||
|
||||
var options = new TelemetryDialOptions { ApiKey = "test-key" };
|
||||
if (mode is not null)
|
||||
{
|
||||
options.Mode = mode;
|
||||
}
|
||||
|
||||
services.AddSingleton<IOptions<TelemetryDialOptions>>(Options.Create(options));
|
||||
|
||||
services.AddAkka("otopcua-test", (ab, _) =>
|
||||
{
|
||||
ab.AddHocon(@"
|
||||
akka.actor.provider = ""Akka.Cluster.ClusterActorRefProvider, Akka.Cluster""
|
||||
akka.remote.dot-netty.tcp.hostname = ""127.0.0.1""
|
||||
akka.remote.dot-netty.tcp.port = 0
|
||||
akka.cluster.seed-nodes = []
|
||||
akka.cluster.roles = [""admin""]
|
||||
", HoconAddMode.Prepend);
|
||||
ab.WithOtOpcUaSignalRBridges();
|
||||
});
|
||||
})
|
||||
.Build();
|
||||
|
||||
/// <summary>An <see cref="IDbContextFactory{TContext}"/> whose contexts share one InMemory database.</summary>
|
||||
private sealed class InMemoryConfigDbFactory(string dbName) : IDbContextFactory<OtOpcUaConfigDbContext>
|
||||
{
|
||||
public OtOpcUaConfigDbContext CreateDbContext() =>
|
||||
new(new DbContextOptionsBuilder<OtOpcUaConfigDbContext>()
|
||||
.UseInMemoryDatabase(dbName)
|
||||
.Options);
|
||||
|
||||
public Task<OtOpcUaConfigDbContext> CreateDbContextAsync(CancellationToken cancellationToken = default) =>
|
||||
Task.FromResult(CreateDbContext());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user