diff --git a/docs/plans/2026-07-21-per-cluster-mesh-phase1.md.tasks.json b/docs/plans/2026-07-21-per-cluster-mesh-phase1.md.tasks.json new file mode 100644 index 00000000..21816038 --- /dev/null +++ b/docs/plans/2026-07-21-per-cluster-mesh-phase1.md.tasks.json @@ -0,0 +1,18 @@ +{ + "plan": "docs/plans/2026-07-21-per-cluster-mesh-phase1.md", + "branch": "feat/mesh-phase1", + "decisions": { + "seal-semantics": "DB-derived + Enabled hatch (confirmed 2026-07-22)", + "cluster-scope": "Dropped — no cluster-scoped deployment exists; expected set = all enabled ClusterNode rows", + "node-roles": "Every ClusterNode row is a driver node; documented + pinned by test, no Roles column", + "adminui-edit": "Deferred to Phase 2 (nothing reads the columns until then)" + }, + "tasks": [ + { "id": 1, "subject": "Address columns on ClusterNode", "status": "pending", "blockedBy": [] }, + { "id": 2, "subject": "EF migration AddClusterNodeTransportPorts", "status": "pending", "blockedBy": [1] }, + { "id": 3, "subject": "Coordinator sources expected-ack set from ClusterNode", "status": "pending", "blockedBy": [2] }, + { "id": 4, "subject": "Reconcile duplicated AkkaPort with the node's own Cluster:Port", "status": "pending", "blockedBy": [2] }, + { "id": 5, "subject": "Rig seed + docs", "status": "pending", "blockedBy": [3, 4] }, + { "id": 6, "subject": "Live gate on docker-dev", "status": "pending", "blockedBy": [5] } + ] +} diff --git a/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/Entities/ClusterNode.cs b/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/Entities/ClusterNode.cs index 91a9bf02..52ecfddd 100644 --- a/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/Entities/ClusterNode.cs +++ b/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/Entities/ClusterNode.cs @@ -18,6 +18,33 @@ public sealed class ClusterNode /// The dashboard HTTP port (default 8081). public int DashboardPort { get; set; } = 8081; + /// + /// Akka remoting port (default 4053). This is central's dial target, not the node's own + /// binding configuration — the node binds from Cluster:Port in its own appsettings, + /// and this column duplicates that value for a reader that cannot see the node's config. + /// Per-cluster mesh Phase 2 builds its ClusterClient contact points from + /// + this port, at which point central and the site node no longer share a + /// gossip ring and central has no other way to learn it. + /// + /// + /// The duplication is real and unenforced by the schema. ClusterNodeAddressReconciler + /// on the admin node compares this row against observed cluster membership and logs an Error + /// on mismatch — a node that binds 4054 while its row says 4053 is unreachable from central in + /// Phase 2, and the symptom there is a silent absence of acks rather than an error. + /// + public int AkkaPort { get; set; } = 4053; + + /// + /// gRPC port central dials for the Phase 5 telemetry stream, or when the + /// node exposes none. Also central's dial target rather than the node's binding config. + /// + /// + /// Nullable by intent: nothing listens on this port yet, and a non-null default would assert a + /// port that does not exist. Phase 5 populates it; until then is the + /// honest value. + /// + public int? GrpcPort { get; set; } + /// /// OPC UA ApplicationUri — MUST be unique per node per OPC UA spec. Clients pin trust here. /// Fleet-wide unique index enforces no two nodes share a value. diff --git a/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/OtOpcUaConfigDbContext.cs b/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/OtOpcUaConfigDbContext.cs index defa00de..6c380a3c 100644 --- a/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/OtOpcUaConfigDbContext.cs +++ b/src/Core/ZB.MOM.WW.OtOpcUa.Configuration/OtOpcUaConfigDbContext.cs @@ -142,6 +142,11 @@ public sealed class OtOpcUaConfigDbContext(DbContextOptions x.Host).HasMaxLength(255); e.Property(x => x.ApplicationUri).HasMaxLength(256); e.Property(x => x.DriverConfigOverridesJson).HasColumnType("nvarchar(max)"); + // Central's dial targets (per-cluster mesh Phase 1). AkkaPort carries a DB-side default so + // rows migrated from before the column existed come out at 4053 rather than 0 — every node + // in the fleet does listen on a remoting port, so 0 is never a truthful value for it. + // GrpcPort has no default: nothing listens on it until Phase 5. + e.Property(x => x.AkkaPort).HasDefaultValue(4053); e.Property(x => x.LastSeenAt).HasColumnType("datetime2(3)"); e.Property(x => x.CreatedAt).HasColumnType("datetime2(3)").HasDefaultValueSql("SYSUTCDATETIME()"); e.Property(x => x.CreatedBy).HasMaxLength(128); diff --git a/tests/Core/ZB.MOM.WW.OtOpcUa.Configuration.Tests/ClusterNodeTransportPortsTests.cs b/tests/Core/ZB.MOM.WW.OtOpcUa.Configuration.Tests/ClusterNodeTransportPortsTests.cs new file mode 100644 index 00000000..cda4bce5 --- /dev/null +++ b/tests/Core/ZB.MOM.WW.OtOpcUa.Configuration.Tests/ClusterNodeTransportPortsTests.cs @@ -0,0 +1,112 @@ +using Microsoft.EntityFrameworkCore; +using Shouldly; +using Xunit; +using ZB.MOM.WW.OtOpcUa.Configuration.Entities; +using ZB.MOM.WW.OtOpcUa.Configuration.Enums; + +namespace ZB.MOM.WW.OtOpcUa.Configuration.Tests; + +/// +/// Round-trips / through +/// the real schema (per-cluster mesh Phase 1). These are central's dial targets: Phase 2 +/// builds its ClusterClient contact points from Host + AkkaPort, and Phase 5 dials +/// GrpcPort for the telemetry stream. Central cannot see a site node's appsettings once +/// the meshes split, so the value has to live in a row central can read. +/// +[Trait("Category", "SchemaCompliance")] +[Collection(nameof(SchemaComplianceCollection))] +public sealed class ClusterNodeTransportPortsTests(SchemaComplianceFixture fixture) +{ + /// Verifies both ports survive a DB round-trip and that AkkaPort defaults to 4053. + [Fact] + public async Task Transport_ports_round_trip_and_AkkaPort_defaults_to_4053() + { + await using var ctx = NewContext(); + var clusterId = await SeedClusterAsync(ctx); + + // Explicit non-default values on one node — proves the columns are mapped and persisted + // rather than silently dropped (an unmapped property round-trips fine in memory). + ctx.ClusterNodes.Add(NewNode(clusterId, "explicit", n => + { + n.AkkaPort = 4055; + n.GrpcPort = 5223; + })); + // Nothing set on the other — AkkaPort must land on 4053 and GrpcPort must stay null. + // GrpcPort is deliberately nullable: nothing listens on it until Phase 5, and a non-null + // default would assert a port that does not exist. + ctx.ClusterNodes.Add(NewNode(clusterId, "defaulted")); + await ctx.SaveChangesAsync(); + + await using var read = NewContext(); + var explicitNode = await read.ClusterNodes.AsNoTracking() + .SingleAsync(n => n.NodeId == $"{clusterId}-explicit:4053"); + explicitNode.AkkaPort.ShouldBe(4055); + explicitNode.GrpcPort.ShouldBe(5223); + + var defaulted = await read.ClusterNodes.AsNoTracking() + .SingleAsync(n => n.NodeId == $"{clusterId}-defaulted:4053"); + defaulted.AkkaPort.ShouldBe(4053); + defaulted.GrpcPort.ShouldBeNull(); + } + + /// + /// Verifies the columns are queryable server-side — Phase 2's contact-point refresh selects + /// them in SQL rather than materialising every node, so a client-side-evaluation regression + /// (or an unmapped property) has to fail here. + /// + [Fact] + public async Task Transport_ports_are_queryable_server_side() + { + await using var ctx = NewContext(); + var clusterId = await SeedClusterAsync(ctx); + ctx.ClusterNodes.Add(NewNode(clusterId, "a", n => n.GrpcPort = 5223)); + ctx.ClusterNodes.Add(NewNode(clusterId, "b")); + await ctx.SaveChangesAsync(); + + await using var read = NewContext(); + var contacts = await read.ClusterNodes.AsNoTracking() + .Where(n => n.ClusterId == clusterId && n.GrpcPort != null) + .Select(n => new { n.Host, n.AkkaPort, n.GrpcPort }) + .ToListAsync(); + + contacts.Count.ShouldBe(1); + contacts[0].AkkaPort.ShouldBe(4053); + contacts[0].GrpcPort.ShouldBe(5223); + } + + private static ClusterNode NewNode(string clusterId, string suffix, Action? configure = null) + { + var node = new ClusterNode + { + NodeId = $"{clusterId}-{suffix}:4053", + ClusterId = clusterId, + Host = $"{clusterId}-{suffix}", + ApplicationUri = $"urn:OtOpcUa:{clusterId}-{suffix}", + CreatedBy = nameof(ClusterNodeTransportPortsTests), + }; + configure?.Invoke(node); + return node; + } + + private static async Task SeedClusterAsync(OtOpcUaConfigDbContext ctx) + { + // Unique per test: the collection shares one database, and ApplicationUri is fleet-wide unique. + var clusterId = $"PORTS-{Guid.NewGuid():N}"[..24]; + ctx.ServerClusters.Add(new ServerCluster + { + ClusterId = clusterId, + Name = clusterId, + Enterprise = "zb", + Site = "ports-test", + RedundancyMode = RedundancyMode.Warm, + CreatedBy = nameof(ClusterNodeTransportPortsTests), + }); + await ctx.SaveChangesAsync(); + return clusterId; + } + + private OtOpcUaConfigDbContext NewContext() => + new(new DbContextOptionsBuilder() + .UseSqlServer(fixture.ConnectionString) + .Options); +}