diff --git a/Darling/Darling.Tests/AgGroupIdLiveReaderTests.cs b/Darling/Darling.Tests/AgGroupIdLiveReaderTests.cs new file mode 100644 index 000000000..2ab153001 --- /dev/null +++ b/Darling/Darling.Tests/AgGroupIdLiveReaderTests.cs @@ -0,0 +1,119 @@ +/* + * Copyright (c) 2026 Erik Darling, Darling Data LLC + * + * This file is part of the SQL Server Performance Monitor. + * + * Licensed under the MIT License. See LICENSE file in the project root for full license information. + */ + +using System; +using System.Threading.Tasks; +using Npgsql; +using PerformanceMonitor.Collectors; +using PerformanceMonitor.Common; +using PerformanceMonitor.Darling.Service.Mcp; +using PerformanceMonitor.Darling.Storage; +using PerformanceMonitor.Darling.Viewer; +using Xunit; + +namespace Darling.Tests; + +/* #1776 own-store: this fact mints its own scratch database through ScratchPostgres and never touches another + test's rows, so it is deliberately NOT [Collection("live-postgres")]. */ + +/// +/// V151/#4475 live reader pin: two monitored SECONDARIES of one real Availability Group, its primary +/// unmonitored, carrying the SAME group_id and DISJOINT replica-name sets, count as ONE group through +/// the real reads -- the MCP/web AG reader's distinct_ag_count and the Viewer's +/// count over its own read -- not just through called directly. +/// +public sealed class AgGroupIdLiveReaderTests +{ + private const string ServerNameA = "darling-ag-group-id-a"; + private const string ServerNameB = "darling-ag-group-id-b"; + private const string GroupId = "11111111-2222-3333-4444-555555555555"; + + private static string? BaseConnectionString => Environment.GetEnvironmentVariable("DARLING_TEST_PG"); + + [Fact] + public async Task DistinctAgCount_TwoSecondariesSameGroupId_DisjointReplicaSets_IsOneThroughTheRealReads() + { + var baseConnectionString = BaseConnectionString; + Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString), "Set DARLING_TEST_PG to run the live AG group_id reader pin."); + + var ct = TestContext.Current.CancellationToken; + var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct); + var bodySucceeded = false; + try + { + await using var connection = new NpgsqlConnection(scratch.ConnectionString); + await connection.OpenAsync(ct); + await PgMigrations.MigrateAsync(connection, ct); + + /* If the branch's V151 rung has not landed yet at the time this runs, plant the column ourselves + so this pin still compiles and runs against a store that predates it -- said explicitly rather + than silently swallowed. */ + await using (var addColumn = new NpgsqlCommand( + "ALTER TABLE ag_replica_states ADD COLUMN IF NOT EXISTS group_id text;", connection)) + { + await addColumn.ExecuteNonQueryAsync(ct); + } + + var serverIdA = ServerIdHelper.GetDeterministicHashCode(ServerNameA); + var serverIdB = ServerIdHelper.GetDeterministicHashCode(ServerNameB); + await DarlingMcpTestData.RegisterServerAsync(connection, serverIdA, ServerNameA, ct); + await DarlingMcpTestData.RegisterServerAsync(connection, serverIdB, ServerNameB, ct); + + var when = DarlingMcpTestData.TruncateToSeconds(DateTime.UtcNow).AddMinutes(-5); + + /* Each server's replica-states row reports ONLY ITSELF -- the real shape a monitored SECONDARY + produces under sys.dm_hadr_availability_replica_states' local-only rule, with the AG's primary + unmonitored. The replica-name sets are DISJOINT ({AGNODE-A} vs {AGNODE-B}); only the shared + group_id links them. */ + await DarlingMcpTestData.ExecAsync(connection, ct, + @"INSERT INTO ag_replica_states (collection_id, collection_time, server_id, server_name, ag_name, replica_server_name, role_desc, operational_state_desc, connected_state_desc, recovery_health_desc, synchronization_health_desc, availability_mode_desc, failover_mode_desc, endpoint_url, group_id) +VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)", + CollectionIdGenerator.Next(), DarlingMcpTestData.Naive(when), serverIdA, ServerNameA, "AG1", "AGNODE-A", "SECONDARY", + "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://AGNODE-A:5022", GroupId); + + await DarlingMcpTestData.ExecAsync(connection, ct, + @"INSERT INTO ag_replica_states (collection_id, collection_time, server_id, server_name, ag_name, replica_server_name, role_desc, operational_state_desc, connected_state_desc, recovery_health_desc, synchronization_health_desc, availability_mode_desc, failover_mode_desc, endpoint_url, group_id) +VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)", + CollectionIdGenerator.Next(), DarlingMcpTestData.Naive(when), serverIdB, ServerNameB, "AG1", "AGNODE-B", "SECONDARY", + "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://AGNODE-B:5022", GroupId); + + /* The MCP/web read (get_ag_health's distinct_ag_count), through DarlingAgReader.GetAgHealthAsync -- + the same code path /api/ag and get_ag_health call, not AgTopology.CountDistinctGroups called + directly. */ + await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString); + var mcpResult = await DarlingAgReader.GetAgHealthAsync(postgres, null, DateTime.UtcNow, cancellationToken: ct); + Assert.Equal(2, mcpResult.AvailabilityGroupCount); + Assert.Equal(1, mcpResult.DistinctAgCount); + + /* The Viewer's own read and count, through ViewerDataService.GetAvailabilityGroupsAsync and + AgTopology.Counts over its cards -- the same call path the WPF AG tab uses, not a manual + reconstruction of its cards. */ + var viewerService = new ViewerDataService(scratch.ConnectionString); + try + { + var cards = await viewerService.GetAvailabilityGroupsAsync(ct); + var (groups, reportingServers, views) = AgTopology.Counts(cards); + + Assert.Equal(2, views); + Assert.Equal(2, reportingServers); + Assert.Equal(1, groups); + } + finally + { + await viewerService.DisposeAsync(); + } + + bodySucceeded = true; + } + finally + { + await LiveStoreCleanup.RunAsync(scratch.ConnectionString, bodySucceeded, async (_, _) => await Task.CompletedTask); + await scratch.DisposeAsync(); + } + } +} diff --git a/Darling/Darling.Tests/AgGroupIdRungTests.cs b/Darling/Darling.Tests/AgGroupIdRungTests.cs new file mode 100644 index 000000000..0f3e501ef --- /dev/null +++ b/Darling/Darling.Tests/AgGroupIdRungTests.cs @@ -0,0 +1,370 @@ +/* + * Copyright (c) 2026 Erik Darling, Darling Data LLC + * + * This file is part of the SQL Server Performance Monitor. + * + * Licensed under the MIT License. See LICENSE file in the project root for full license information. + */ + +using System; +using System.Collections.Generic; +using System.Globalization; +using System.Linq; +using System.Threading.Tasks; +using Npgsql; +using PerformanceMonitor.Darling.Storage; +using PerformanceMonitor.Darling.Viewer; +using Xunit; + +namespace Darling.Tests; + +/// +/// Pins Darling rung V151 (#4475): group_id, a text column on both collect.ag_replica_states +/// and collect.ag_database_replica_states that carries sys.availability_groups.group_id, the +/// same GUID (stored as a 36-character string) on every replica of one physical Availability Group. This +/// file is the RUNG (ladder, viewer probe) and the schema-after-migrate proof: both columns exist on a +/// fresh migrate, the ALTER succeeds on a store whose chunks are already compressed, a pre-V151 row reads +/// group_id IS NULL rather than erroring, and an INSERT in the collector's own column order round +/// trips a GUID string end to end. +/// +/// This file's "I am the top rung" claim takes over from +/// CollectionLogWatermarkAndJobHistoryIndexesRungTests (V150) now that V151 has landed. +/// +/* #1776 own-store: each fact mints its own scratch database through ScratchPostgres and never touches the + shared store's tables, so it cannot race the live collection and serializing it would be pure slowdown. */ +public sealed class AgGroupIdRungTests +{ + private const int RungVersion = 151; + private const int PreviousVersion = 150; + + /// This rung's sentinel ordinal in the viewer probe — the newest, so the last argument. + private const int ProbeOrdinal = 126; + + private static string? ConnectionString => Environment.GetEnvironmentVariable("DARLING_TEST_PG"); + + /// + /// The rung is registered and is the new top of the ladder — the claim this class takes over from + /// CollectionLogWatermarkAndJobHistoryIndexesRungTests (V150) now that V151 has landed. + /// + [Fact] + public void TheRungIsRegisteredAtTheTopOfADenseLadder() + { + var versions = PgMigrations.Scripts.Select(s => s.Version).ToList(); + + Assert.Equal("ag-group-id", PgMigrations.Scripts.Single(s => s.Version == RungVersion).Name); + Assert.Equal(StorageVersion.SchemaVersion, PgMigrations.Scripts[^1].Version); + Assert.Equal(StorageVersion.SchemaVersion, versions.Max()); + Assert.Equal(RungVersion, StorageVersion.SchemaVersion); + Assert.Equal(versions.Distinct().OrderBy(v => v), versions); + } + + /// + /// The viewer probe's sentinel carries this rung, and the map treats it as the TOP arm: a missing top arm + /// maps a fully-migrated store one rung short, permanently, because + /// is . + /// + [Fact] + public void TheProbeMapsAFullyMigratedStoreToThisTopRung() + { + var probe = ViewerDataService.StoreSchemaProbeSql.Replace("\r\n", "\n", StringComparison.Ordinal); + Assert.Contains("table_name = 'ag_replica_states'", probe, StringComparison.Ordinal); + Assert.Contains("column_name = 'group_id'", probe, StringComparison.Ordinal); + + var viewer = RepoFile.ReadRepoFile("Darling", "PerformanceMonitor.Darling.Viewer", "ViewerDataService.cs"); + Assert.Contains($"reader.GetBoolean({ProbeOrdinal})", viewer, StringComparison.Ordinal); + Assert.DoesNotContain($"reader.GetBoolean({ProbeOrdinal + 1})", viewer, StringComparison.Ordinal); + + Assert.Equal(StorageVersion.SchemaVersion, ViewerDataService.RequiredStoreSchemaVersion); + + var method = typeof(ViewerDataService).GetMethod("MapProbedSchemaVersion", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Static)!; + var arity = method.GetParameters().Length; + Assert.Equal(ProbeOrdinal, arity - 1); + Assert.Equal("hasAgGroupId", method.GetParameters()[ProbeOrdinal].Name); + + var all = Enumerable.Repeat((object)true, arity).ToArray(); + Assert.Equal(StorageVersion.SchemaVersion, (int)method.Invoke(null, all)!); + + var behind = (object[])all.Clone(); + behind[ProbeOrdinal] = false; + Assert.Equal(PreviousVersion, (int)method.Invoke(null, behind)!); + + var thisArm = viewer.IndexOf("if (hasAgGroupId)", StringComparison.Ordinal); + var previousArm = viewer.IndexOf("if (hasCollectionLogWatermarkAndJobHistoryIndexes)", StringComparison.Ordinal); + Assert.True(thisArm >= 0, "the viewer has no V151 sentinel arm — a fully-migrated store would map one rung short"); + Assert.True(thisArm < previousArm, "the V151 arm sits below V150's, so a current store maps one rung short"); + Assert.Contains( + "return " + StorageVersion.SchemaVersion.ToString(CultureInfo.InvariantCulture) + ";", + viewer[thisArm..previousArm], StringComparison.Ordinal); + } + + /// + /// The LIVE schema after migrate: both columns exist. Run against origin/dev (pre-V151) this is + /// RED — neither column exists — proving the pin actually checks the rung rather than a tautology. + /// + [Fact] + public async Task AfterMigrate_BothGroupIdColumnsExist() + { + var baseConnectionString = ConnectionString; + Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString), + "Set DARLING_TEST_PG to a Postgres connection string to run the V151 schema pin."); + + var ct = TestContext.Current.CancellationToken; + + await using var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct); + await using var connection = new NpgsqlConnection(scratch.ConnectionString); + await connection.OpenAsync(ct); + await PgMigrations.MigrateAsync(connection, ct); + + Assert.True(await ColumnExistsAsync(connection, "ag_replica_states", "group_id", ct), + "V151 adds ag_replica_states.group_id"); + Assert.True(await ColumnExistsAsync(connection, "ag_database_replica_states", "group_id", ct), + "V151 adds ag_database_replica_states.group_id"); + } + + /// + /// Rerunning the migration ladder (as startup does on an already-migrated store) is idempotent: the + /// two ADD COLUMN IF NOT EXISTS statements do not error and both columns still exist. + /// + [Fact] + public async Task MigrateAsync_RunTwice_IsIdempotent() + { + var baseConnectionString = ConnectionString; + Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString), + "Set DARLING_TEST_PG to a Postgres connection string to run the V151 idempotency pin."); + + var ct = TestContext.Current.CancellationToken; + + await using var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct); + await using var connection = new NpgsqlConnection(scratch.ConnectionString); + await connection.OpenAsync(ct); + await PgMigrations.MigrateAsync(connection, ct); + await PgMigrations.MigrateAsync(connection, ct); + + Assert.True(await ColumnExistsAsync(connection, "ag_replica_states", "group_id", ct)); + Assert.True(await ColumnExistsAsync(connection, "ag_database_replica_states", "group_id", ct)); + } + + /// + /// The live migration round trip on COMPRESSED hypertables: climb to V150, plant one row in each of + /// ag_replica_states and ag_database_replica_states, compress their chunks the product + /// way (create_hypertable + timescaledb.compress + compress_chunk — the same + /// three statements and + /// issue at startup), then + /// roll back JUST V151's version stamp and re-apply it. The ALTER must succeed against a hypertable + /// whose chunks are already compressed — a TimescaleDB compressed chunk rejects some DDL forms outright + /// — and the planted pre-V151 rows must read group_id IS NULL rather than erroring or losing the + /// row. A closing INSERT in the collector's own column order ('s + /// WritePayload order, group_id last) carries a GUID string and reads back unchanged. + /// + [Fact] + public async Task AStoreAtV150WithCompressedChunks_MigratesToV151_AndOldRowsReadGroupIdNull() + { + var baseConnectionString = ConnectionString; + Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString), + "Set DARLING_TEST_PG to a Postgres connection string to run the V151 compressed-hypertable round trip."); + + var ct = TestContext.Current.CancellationToken; + + await using var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct); + await using var connection = new NpgsqlConnection(scratch.ConnectionString); + await connection.OpenAsync(ct); + + /* Climb all the way to the top first — the compressed-hypertable state below needs the rest of + the schema (servers, config) present, and a fresh scratch database has none of it yet. */ + await PgMigrations.MigrateAsync(connection, ct); + + var oldTime = DateTime.SpecifyKind(new DateTime(2026, 1, 1, 0, 0, 0), DateTimeKind.Unspecified); + + await InsertReplicaRowAsync(connection, ct, collectionId: 1, collectionTime: oldTime, groupIdColumn: false); + await InsertDatabaseReplicaRowAsync(connection, ct, collectionId: 1, collectionTime: oldTime, groupIdColumn: false); + + /* This test MUST run on CI (CI has TimescaleDB); the compressed-hypertable state is the whole + point. A fresh ScratchPostgres database does not inherit the extension, so enable it here on the + scratch connection before the create_hypertable calls below, the same way + DarlingWatermarkFloorScanBoundLiveTests / CaptureDownChunkOrderTests do. */ + var timescaleEnabled = await LiveTimescaleProbe.TryEnableAsync(scratch.ConnectionString, ct); + Assert.True(timescaleEnabled, "TimescaleDB must be available on CI for the compressed-hypertable round trip"); + + /* Convert both tables to hypertables and compress their one chunk each, the product's own way — + this is what "V150's ALTER succeeded on the compressed hypertable" actually needs to test + against, rather than an ordinary heap table. */ + await ExecAsync(connection, ct, + "SELECT create_hypertable('collect.ag_replica_states', by_range('collection_time', INTERVAL '1 days'), if_not_exists => true, migrate_data => true)"); + await ExecAsync(connection, ct, + "SELECT create_hypertable('collect.ag_database_replica_states', by_range('collection_time', INTERVAL '1 days'), if_not_exists => true, migrate_data => true)"); + await ExecAsync(connection, ct, + "ALTER TABLE collect.ag_replica_states SET (timescaledb.compress, timescaledb.compress_segmentby = 'server_id')"); + await ExecAsync(connection, ct, + "ALTER TABLE collect.ag_database_replica_states SET (timescaledb.compress, timescaledb.compress_segmentby = 'server_id')"); + await CompressAllChunksAsync(connection, ct, "collect.ag_replica_states"); + await CompressAllChunksAsync(connection, ct, "collect.ag_database_replica_states"); + + Assert.True(await AnyChunkCompressedAsync(connection, ct, "collect.ag_replica_states"), + "the seed row's chunk did not compress — the ALTER below would prove nothing about a compressed hypertable"); + Assert.True(await AnyChunkCompressedAsync(connection, ct, "collect.ag_database_replica_states")); + + /* Roll back JUST V151's version stamp, so the store looks exactly like a V150 store whose chunks + are already compressed. */ + await ExecAsync(connection, ct, "DELETE FROM darling_schema_version WHERE version >= " + RungVersion); + + using (var version = new NpgsqlCommand("SELECT MAX(version) FROM darling_schema_version", connection)) + { + Assert.Equal(PreviousVersion, Convert.ToInt32(await version.ExecuteScalarAsync(ct))); + } + + /* The upgrade: apply every rung above V150, which is exactly V151 on this build. */ + var applied = await PgMigrations.MigrateAsync(connection, ct); + Assert.Equal(PgMigrations.Scripts.Count(m => m.Version >= RungVersion), applied); + + Assert.True(await ColumnExistsAsync(connection, "ag_replica_states", "group_id", ct)); + Assert.True(await ColumnExistsAsync(connection, "ag_database_replica_states", "group_id", ct)); + + Assert.Null(await ReadGroupIdAsync(connection, ct, "ag_replica_states", collectionId: 1)); + Assert.Null(await ReadGroupIdAsync(connection, ct, "ag_database_replica_states", collectionId: 1)); + + /* A closing INSERT in the collector's own column order (group_id last), carrying a real GUID + string, round trips end to end on the now-widened, compressed hypertable. */ + var newTime = DateTime.SpecifyKind(new DateTime(2026, 9, 27, 0, 0, 0), DateTimeKind.Unspecified); + const string groupId = "3F2504E0-4F89-11D3-9A0C-0305E82C3301"; + + await InsertReplicaRowAsync(connection, ct, collectionId: 2, collectionTime: newTime, groupIdColumn: true, groupId); + await InsertDatabaseReplicaRowAsync(connection, ct, collectionId: 2, collectionTime: newTime, groupIdColumn: true, groupId); + + Assert.Equal(groupId, await ReadGroupIdAsync(connection, ct, "ag_replica_states", collectionId: 2)); + Assert.Equal(groupId, await ReadGroupIdAsync(connection, ct, "ag_database_replica_states", collectionId: 2)); + } + + private static async Task InsertReplicaRowAsync( + NpgsqlConnection connection, System.Threading.CancellationToken ct, long collectionId, DateTime collectionTime, + bool groupIdColumn, string? groupId = null) + { + var columns = "collection_id, collection_time, server_id, server_name, ag_name, replica_server_name, role_desc, " + + "operational_state_desc, connected_state_desc, recovery_health_desc, synchronization_health_desc, " + + "availability_mode_desc, failover_mode_desc, endpoint_url, is_local" + (groupIdColumn ? ", group_id" : string.Empty); + var placeholders = groupIdColumn + ? "$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16" + : "$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15"; + + await using var command = new NpgsqlCommand( + $"INSERT INTO collect.ag_replica_states ({columns}) VALUES ({placeholders})", connection); + command.Parameters.AddWithValue(collectionId); + command.Parameters.AddWithValue(collectionTime); + command.Parameters.AddWithValue(1); + command.Parameters.AddWithValue("server1"); + command.Parameters.AddWithValue("AG1"); + command.Parameters.AddWithValue("NODE1"); + command.Parameters.AddWithValue("PRIMARY"); + command.Parameters.AddWithValue("ONLINE"); + command.Parameters.AddWithValue("CONNECTED"); + command.Parameters.AddWithValue("ONLINE"); + command.Parameters.AddWithValue("HEALTHY"); + command.Parameters.AddWithValue("SYNCHRONOUS_COMMIT"); + command.Parameters.AddWithValue("AUTOMATIC"); + command.Parameters.AddWithValue("TCP://NODE1:5022"); + command.Parameters.AddWithValue(true); + if (groupIdColumn) + { + command.Parameters.AddWithValue(groupId!); + } + + await command.ExecuteNonQueryAsync(ct); + } + + private static async Task InsertDatabaseReplicaRowAsync( + NpgsqlConnection connection, System.Threading.CancellationToken ct, long collectionId, DateTime collectionTime, + bool groupIdColumn, string? groupId = null) + { + var columns = "collection_id, collection_time, server_id, server_name, ag_name, database_name, replica_server_name, " + + "is_local, synchronization_state_desc, last_hardened_lsn, last_commit_lsn, log_send_queue_size, " + + "redo_queue_size, log_send_rate, redo_rate, is_suspended, suspend_reason_desc, availability_mode_desc, " + + "secondary_lag_seconds" + (groupIdColumn ? ", group_id" : string.Empty); + var placeholders = groupIdColumn + ? "$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20" + : "$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19"; + + await using var command = new NpgsqlCommand( + $"INSERT INTO collect.ag_database_replica_states ({columns}) VALUES ({placeholders})", connection); + command.Parameters.AddWithValue(collectionId); + command.Parameters.AddWithValue(collectionTime); + command.Parameters.AddWithValue(1); + command.Parameters.AddWithValue("server1"); + command.Parameters.AddWithValue("AG1"); + command.Parameters.AddWithValue("db1"); + command.Parameters.AddWithValue("NODE1"); + command.Parameters.AddWithValue(true); + command.Parameters.AddWithValue("SYNCHRONIZED"); + command.Parameters.AddWithValue("0/0"); + command.Parameters.AddWithValue("0/0"); + command.Parameters.AddWithValue(0L); + command.Parameters.AddWithValue(0L); + command.Parameters.AddWithValue(0L); + command.Parameters.AddWithValue(0L); + command.Parameters.AddWithValue(false); + command.Parameters.AddWithValue(DBNull.Value); + command.Parameters.AddWithValue("SYNCHRONOUS_COMMIT"); + command.Parameters.AddWithValue(0L); + if (groupIdColumn) + { + command.Parameters.AddWithValue(groupId!); + } + + await command.ExecuteNonQueryAsync(ct); + } + + private static async Task ReadGroupIdAsync( + NpgsqlConnection connection, System.Threading.CancellationToken ct, string table, long collectionId) + { + await using var command = new NpgsqlCommand( + $"SELECT group_id FROM collect.{table} WHERE collection_id = $1", connection); + command.Parameters.AddWithValue(collectionId); + var value = await command.ExecuteScalarAsync(ct); + return value is DBNull or null ? null : (string)value; + } + + private static async Task ExecAsync(NpgsqlConnection connection, System.Threading.CancellationToken ct, string sql) + { + await using var command = new NpgsqlCommand(sql, connection); + await command.ExecuteNonQueryAsync(ct); + } + + private static async Task ColumnExistsAsync( + NpgsqlConnection connection, string table, string column, System.Threading.CancellationToken ct) + { + await using var command = new NpgsqlCommand( + "SELECT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema = 'collect' AND table_name = $1 AND column_name = $2)", + connection); + command.Parameters.AddWithValue(table); + command.Parameters.AddWithValue(column); + return (bool)(await command.ExecuteScalarAsync(ct))!; + } + + private static async Task CompressAllChunksAsync( + NpgsqlConnection connection, System.Threading.CancellationToken ct, string table) + { + var chunks = new List(); + await using (var chunkList = new NpgsqlCommand($"SELECT show_chunks('{table}')::text", connection)) + await using (var reader = await chunkList.ExecuteReaderAsync(ct)) + { + while (await reader.ReadAsync(ct)) + { + chunks.Add(reader.GetString(0)); + } + } + + foreach (var chunk in chunks) + { + await using var compress = new NpgsqlCommand($"SELECT compress_chunk('{chunk}', if_not_compressed => true)", connection); + await compress.ExecuteNonQueryAsync(ct); + } + } + + private static async Task AnyChunkCompressedAsync( + NpgsqlConnection connection, System.Threading.CancellationToken ct, string table) + { + await using var command = new NpgsqlCommand( + "SELECT EXISTS (SELECT 1 FROM timescaledb_information.chunks WHERE hypertable_name = " + + "split_part($1, '.', 2) AND is_compressed)", connection); + command.Parameters.AddWithValue(table); + return (bool)(await command.ExecuteScalarAsync(ct))!; + } +} diff --git a/Darling/Darling.Tests/CollectionLogWatermarkAndJobHistoryIndexesRungTests.cs b/Darling/Darling.Tests/CollectionLogWatermarkAndJobHistoryIndexesRungTests.cs index b98414d6e..e8f910c5e 100644 --- a/Darling/Darling.Tests/CollectionLogWatermarkAndJobHistoryIndexesRungTests.cs +++ b/Darling/Darling.Tests/CollectionLogWatermarkAndJobHistoryIndexesRungTests.cs @@ -24,8 +24,9 @@ namespace Darling.Tests; /// (ladder, viewer probe) and the schema-after-migrate proof: both indexes exist, plain /// CREATE INDEX IF NOT EXISTS, idempotent on rerun. /// -/// This file's "I am the top rung" claim takes over from QueryStoreLivenessHotTouchLiveTests -/// (V149) now that V150 has landed. +/// This file's "I am the top rung" claim moved to AgGroupIdRungTests (V151) now that V151 +/// has landed; this file's own rung/probe facts below keep asserting what stays true forever (present, +/// in-order, gated behind the arm above it) rather than "is exactly the top". /// /* #1776 own-store: each fact mints its own scratch database through ScratchPostgres and never touches the shared store's tables, so it cannot race the live collection and serializing it would be pure slowdown. */ @@ -34,34 +35,36 @@ public sealed class CollectionLogWatermarkAndJobHistoryIndexesRungTests private const int RungVersion = 150; private const int PreviousVersion = 149; - /// This rung's sentinel ordinal in the viewer probe — the newest, so the last argument. + /// This rung's sentinel ordinal in the viewer probe — no longer the newest, since V151 landed + /// above it. private const int ProbeOrdinal = 125; private static string? ConnectionString => Environment.GetEnvironmentVariable("DARLING_TEST_PG"); /// - /// The rung is registered and is the new top of the ladder — the claim this class takes over from - /// QueryStoreLivenessHotTouchLiveTests (V149) now that V150 has landed. + /// The rung is registered, and the ladder stays dense above the historical gap — the claim this class + /// took over from QueryStoreLivenessHotTouchLiveTests (V149) moved on again to + /// AgGroupIdRungTests (V151) now that V151 has landed. /// [Fact] - public void TheRungIsRegisteredAtTheTopOfADenseLadder() + public void TheRungIsRegistered_AndTheLadderIsDenseAboveTheHistoricalGap() { var versions = PgMigrations.Scripts.Select(s => s.Version).ToList(); Assert.Equal("collection-log-watermark-and-job-history-indexes", PgMigrations.Scripts.Single(s => s.Version == RungVersion).Name); - Assert.Equal(StorageVersion.SchemaVersion, PgMigrations.Scripts[^1].Version); - Assert.Equal(StorageVersion.SchemaVersion, versions.Max()); - Assert.Equal(RungVersion, StorageVersion.SchemaVersion); Assert.Equal(versions.Distinct().OrderBy(v => v), versions); + + var above = versions.Where(v => v > 45).OrderBy(v => v).ToList(); + Assert.Equal(Enumerable.Range(above[0], above.Count), above); } /// - /// The viewer probe's sentinel carries this rung, and the map treats it as the TOP arm: a missing top arm - /// maps a fully-migrated store one rung short, permanently, because + /// The viewer probe's sentinel carries this rung, and the map treats it as an arm gated below the + /// current top's arm — a missing arm maps a fully-migrated store one rung short, permanently, because /// is . /// [Fact] - public void TheProbeMapsAFullyMigratedStoreToThisTopRung() + public void TheProbeCarriesThisRungsSentinel_AndTheArmSitsBelowTheCurrentTop() { var probe = ViewerDataService.StoreSchemaProbeSql.Replace("\r\n", "\n", StringComparison.Ordinal); Assert.Contains("idx_collection_log_watermark", probe, StringComparison.Ordinal); @@ -69,29 +72,31 @@ public void TheProbeMapsAFullyMigratedStoreToThisTopRung() var viewer = RepoFile.ReadRepoFile("Darling", "PerformanceMonitor.Darling.Viewer", "ViewerDataService.cs"); Assert.Contains($"reader.GetBoolean({ProbeOrdinal})", viewer, StringComparison.Ordinal); - Assert.DoesNotContain($"reader.GetBoolean({ProbeOrdinal + 1})", viewer, StringComparison.Ordinal); - - Assert.Equal(StorageVersion.SchemaVersion, ViewerDataService.RequiredStoreSchemaVersion); var method = typeof(ViewerDataService).GetMethod("MapProbedSchemaVersion", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Static)!; var arity = method.GetParameters().Length; - Assert.Equal(ProbeOrdinal, arity - 1); Assert.Equal("hasCollectionLogWatermarkAndJobHistoryIndexes", method.GetParameters()[ProbeOrdinal].Name); + /* Every rung above this one (V151's hasAgGroupId) must also be false, or the map finds the newer + arm first and this assertion is checking the wrong rung's fallthrough. */ var all = Enumerable.Repeat((object)true, arity).ToArray(); - Assert.Equal(StorageVersion.SchemaVersion, (int)method.Invoke(null, all)!); - var behind = (object[])all.Clone(); - behind[ProbeOrdinal] = false; + for (var i = ProbeOrdinal; i < arity; i++) + { + behind[i] = false; + } Assert.Equal(PreviousVersion, (int)method.Invoke(null, behind)!); + /* V151 (#4475) is now the top rung, so this arm no longer needs to be the LAST one — it only has to + sit below the current top's arm, which is what the ladder-dense invariant above already + guarantees is registered ahead of it. */ var thisArm = viewer.IndexOf("if (hasCollectionLogWatermarkAndJobHistoryIndexes)", StringComparison.Ordinal); - var previousArm = viewer.IndexOf("if (hasHotLivenessTouch)", StringComparison.Ordinal); + var topArm = viewer.IndexOf("if (hasAgGroupId)", StringComparison.Ordinal); Assert.True(thisArm >= 0, "the viewer has no V150 sentinel arm — a fully-migrated store would map one rung short"); - Assert.True(thisArm < previousArm, "the V150 arm sits below V149's, so a current store maps one rung short"); + Assert.True(topArm >= 0 && topArm < thisArm, "the current top rung's arm must sit above the V150 arm"); Assert.Contains( - "return " + StorageVersion.SchemaVersion.ToString(System.Globalization.CultureInfo.InvariantCulture) + ";", - viewer[thisArm..previousArm], StringComparison.Ordinal); + "return " + RungVersion.ToString(System.Globalization.CultureInfo.InvariantCulture) + ";", + viewer[thisArm..(viewer.IndexOf("if (hasHotLivenessTouch)", StringComparison.Ordinal))], StringComparison.Ordinal); } /// diff --git a/Darling/Darling.Tests/DarlingAgReaderTests.cs b/Darling/Darling.Tests/DarlingAgReaderTests.cs index 4ca64dbb8..d171bab60 100644 --- a/Darling/Darling.Tests/DarlingAgReaderTests.cs +++ b/Darling/Darling.Tests/DarlingAgReaderTests.cs @@ -209,8 +209,9 @@ private static Reader.ReplicaRow Replica( string? connected = "CONNECTED", string? operational = "ONLINE", string? recoveryHealth = "ONLINE", - bool? isLocal = null) => - new(serverId, serverName, At(1), agName, replicaName, role, isLocal, operational, connected, recoveryHealth, syncHealth, "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://" + replicaName + ":5022"); + bool? isLocal = null, + string? groupId = null) => + new(serverId, serverName, At(1), agName, replicaName, role, isLocal, operational, connected, recoveryHealth, syncHealth, "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://" + replicaName + ":5022", groupId); private static Reader.DatabaseRow Database( int serverId, @@ -221,8 +222,9 @@ private static Reader.DatabaseRow Database( string state = "SYNCHRONIZED", bool isSuspended = false, long? lagSeconds = 0, - int minutesAgo = 1) => - new(serverId, serverName, At(minutesAgo), agName, databaseName, replicaName, true, state, "0x00", "0x00", 0, 0, 1024, 1024, isSuspended, isSuspended ? "USER_ACTION" : null, "SYNCHRONOUS_COMMIT", lagSeconds); + int minutesAgo = 1, + string? groupId = null) => + new(serverId, serverName, At(minutesAgo), agName, databaseName, replicaName, true, state, "0x00", "0x00", 0, 0, 1024, 1024, isSuspended, isSuspended ? "USER_ACTION" : null, "SYNCHRONOUS_COMMIT", lagSeconds, groupId); [Fact] public void Build_OneAgSeenFromTwoServers_StaysTwoGroupsEachNamingItsReporter() @@ -457,4 +459,42 @@ public void DistinctAgCount_SameAgSeenFromItsPrimaryAndItsSecondary_IsOne() Assert.Equal(1, result.DistinctAgCount); Assert.Equal(2, result.AvailabilityGroupCount); } + + [Fact] + public void DistinctAgCount_TwoSecondariesOfOneAgSameGroupId_DisjointReplicaSets_IsOne() + { + /* (a), through this reader's own public entry point rather than calling AgTopology.CountDistinctGroups + directly: two monitored SECONDARIES of one real AG, its primary unmonitored, each reporting only + itself (so their replica sets are disjoint), carrying the SAME group_id. Without the id, the + name+overlap rule alone counts 2 (see DistinctAgCount_SameAgSeenFromItsPrimaryAndItsSecondary_IsOne's + companion below for the id-less shape) -- with it, they collapse to 1. */ + var replicas = new[] + { + Replica(1, "NODE1", "AG1", "NODE1", "SECONDARY", groupId: "GROUP-GUID-1"), + Replica(2, "NODE2", "AG1", "NODE2", "SECONDARY", groupId: "GROUP-GUID-1"), + }; + + var result = Reader.Build(replicas, Array.Empty(), At(0)); + + Assert.Equal(1, result.DistinctAgCount); + Assert.Equal(2, result.AvailabilityGroupCount); + } + + [Fact] + public void DistinctAgCount_TwoSecondariesOfOneAg_NoGroupId_DisjointReplicaSets_IsTwo() + { + /* The id-LESS companion to the pin above, proving the RED this branch closes: the SAME two disjoint- + replica-set secondaries, with NO group_id on either row, count as 2 under the name+overlap rule + alone -- exactly the gap #4475's follow-up (V151) exists to close. */ + var replicas = new[] + { + Replica(1, "NODE1", "AG1", "NODE1", "SECONDARY"), + Replica(2, "NODE2", "AG1", "NODE2", "SECONDARY"), + }; + + var result = Reader.Build(replicas, Array.Empty(), At(0)); + + Assert.Equal(2, result.DistinctAgCount); + Assert.Equal(2, result.AvailabilityGroupCount); + } } diff --git a/Darling/Darling.Tests/PgSchemaGeneratorTests.cs b/Darling/Darling.Tests/PgSchemaGeneratorTests.cs index f4c182781..9542d58be 100644 --- a/Darling/Darling.Tests/PgSchemaGeneratorTests.cs +++ b/Darling/Darling.Tests/PgSchemaGeneratorTests.cs @@ -479,11 +479,12 @@ GenerateFullSchema. Same contract as V24/V25 — and it is the ONLY thing standi upgraded store's physical shape differ from a fresh one. */ var v34 = Lf(PgMigrations.Scripts.Single(m => m.Version == 34).Sql); - /* The REPLICA-grain table is now the same two-migration story as the database grain below: V34 - created its first 10 payload columns and V37 (#1696) appended is_local, so an upgraded store's - shape is V34 + V37 and only their sum equals the generator's current output. */ + /* The REPLICA-grain table is now a three-migration story: V34 created its first 10 payload + columns, V37 (#1696) appended is_local, and V151 (#4475) appended group_id — so an upgraded + store's shape is V34 + V37 + V151 and only their sum equals the generator's current output. */ var replicaColumns = AgReplicaStatesCollector.Instance.PayloadColumns; const int V34ReplicaColumnCount = 10; + const int V37ReplicaColumnCount = 11; Assert.Contains( CollectQualified(new TruncatedSchema(AgReplicaStatesCollector.Instance, V34ReplicaColumnCount)), @@ -493,7 +494,7 @@ shape is V34 + V37 and only their sum equals the generator's current output. */ var v37 = Lf(PgMigrations.Scripts.Single(m => m.Version == 37).Sql); - foreach (var column in replicaColumns.Skip(V34ReplicaColumnCount)) + foreach (var column in replicaColumns.Skip(V34ReplicaColumnCount).Take(V37ReplicaColumnCount - V34ReplicaColumnCount)) { var generatedType = Lf(PgSchemaGenerator.CreateTable(new TruncatedSchema(AgReplicaStatesCollector.Instance, replicaColumns.Count))) .Split('\n') @@ -504,6 +505,21 @@ shape is V34 + V37 and only their sum equals the generator's current output. */ Assert.Contains($"ADD COLUMN IF NOT EXISTS {generatedType}", v37, StringComparison.Ordinal); } + /* V151 (#4475) appends group_id, the last replica-grain column, as its own additive rung — same + contract as V37 above. */ + var v151ReplicaAlter = Lf(PgMigrations.Scripts.Single(m => m.Version == 151).Sql); + + foreach (var column in replicaColumns.Skip(V37ReplicaColumnCount)) + { + var generatedType = Lf(PgSchemaGenerator.CreateTable(new TruncatedSchema(AgReplicaStatesCollector.Instance, replicaColumns.Count))) + .Split('\n') + .Single(l => l.TrimStart().StartsWith(column.Name + " ", StringComparison.Ordinal)) + .Trim() + .TrimEnd(','); + + Assert.Contains($"ADD COLUMN IF NOT EXISTS {generatedType}", v151ReplicaAlter, StringComparison.Ordinal); + } + /* No "V34 was not widened in place" sweep for this grain, unlike the database one below: the TruncatedSchema assertion above already matches V34's ag_replica_states block EXACTLY, which is a strictly stronger statement than any name-absence check. A substring sweep would also be wrong @@ -511,14 +527,16 @@ strictly stronger statement than any name-absence check. A substring sweep would for it finds the other table's legitimate column and fails. */ Assert.Contains("CREATE INDEX IF NOT EXISTS idx_ag_database_replica_states_time ON collect.ag_database_replica_states(server_id, collection_time);", v34, StringComparison.Ordinal); - /* The database-grain table is the one case where a single migration is NOT the whole story: V34 - created its first 15 payload columns and V36 (#991 addendum) appended 6 more, so an upgraded - store's shape is V34 + V36 and only their SUM can equal the generator's current output. + /* The database-grain table is a three-migration story: V34 created its first 15 payload columns, + V36 (#991 addendum) appended 6 more, and V151 (#4475) appended group_id — so an upgraded store's + shape is V34 + V36 + V151 and only their SUM can equal the generator's current output. Reconstruct that here rather than weakening the pin to name-presence — generate the historical - 15-column shape and assert V34 matches it exactly, then assert V36 appends the remaining columns - in order with the generator's own types. Together those two prove fresh == upgraded. */ + 15-column shape and assert V34 matches it exactly, then assert V36 and V151 each append their + own remaining columns in order with the generator's own types. Together those prove fresh == + upgraded. */ var currentColumns = AgDatabaseReplicaStatesCollector.Instance.PayloadColumns; const int V34ColumnCount = 15; + const int V36ColumnCount = 21; Assert.Contains( CollectQualified(new TruncatedSchema(AgDatabaseReplicaStatesCollector.Instance, V34ColumnCount)), @@ -527,7 +545,7 @@ Reconstruct that here rather than weakening the pin to name-presence — generat var v36 = Lf(PgMigrations.Scripts.Single(m => m.Version == 36).Sql); - foreach (var column in currentColumns.Skip(V34ColumnCount)) + foreach (var column in currentColumns.Skip(V34ColumnCount).Take(V36ColumnCount - V34ColumnCount)) { var generatedType = Lf(PgSchemaGenerator.CreateTable(new TruncatedSchema(AgDatabaseReplicaStatesCollector.Instance, currentColumns.Count))) .Split('\n') @@ -538,6 +556,21 @@ Reconstruct that here rather than weakening the pin to name-presence — generat Assert.Contains($"ADD COLUMN IF NOT EXISTS {generatedType}", v36, StringComparison.Ordinal); } + /* V151 (#4475) appends group_id, the last database-grain column, as its own additive rung — same + contract as V36 above. */ + var v151DatabaseAlter = Lf(PgMigrations.Scripts.Single(m => m.Version == 151).Sql); + + foreach (var column in currentColumns.Skip(V36ColumnCount)) + { + var generatedType = Lf(PgSchemaGenerator.CreateTable(new TruncatedSchema(AgDatabaseReplicaStatesCollector.Instance, currentColumns.Count))) + .Split('\n') + .Single(l => l.TrimStart().StartsWith(column.Name + " ", StringComparison.Ordinal)) + .Trim() + .TrimEnd(','); + + Assert.Contains($"ADD COLUMN IF NOT EXISTS {generatedType}", v151DatabaseAlter, StringComparison.Ordinal); + } + /* And V34 must NOT have been widened in place: its CREATE TABLE IF NOT EXISTS is a no-op on a store that already ran it, so editing V34 instead of adding V36 would leave every migrated store short the new columns while fresh installs got them. */ diff --git a/Darling/Darling.Tests/QsCaptureModeRouteKnobToastRungTests.cs b/Darling/Darling.Tests/QsCaptureModeRouteKnobToastRungTests.cs index a3ba1175e..47b949343 100644 --- a/Darling/Darling.Tests/QsCaptureModeRouteKnobToastRungTests.cs +++ b/Darling/Darling.Tests/QsCaptureModeRouteKnobToastRungTests.cs @@ -381,7 +381,14 @@ public void LiteTwinsOnlyTheCaptureModes_AtSchemaV64() Assert.Equal(CollectorTargetEngine.SqlServer, QueryStoreHealthCollector.Instance.TargetEngine); var initializer = RepoFile.ReadRepoFile("Lite", "Database", "DuckDbInitializer.cs"); - Assert.Contains("internal const int CurrentSchemaVersion = 64;", initializer, StringComparison.Ordinal); + /* Stays-true shape (#4475 raised CurrentSchemaVersion past 64, so a literal "= 64;" pin would + break on the next rung): the v64 block must still be present, and CurrentSchemaVersion must be + AT LEAST 64 — this fact is about Lite twinning ONLY the capture modes at v64, not about v64 + being the ceiling. */ + Assert.Contains("internal const int CurrentSchemaVersion = ", initializer, StringComparison.Ordinal); + var versionMatch = System.Text.RegularExpressions.Regex.Match(initializer, @"internal const int CurrentSchemaVersion = (\d+);"); + Assert.True(versionMatch.Success, "DuckDbInitializer has no CurrentSchemaVersion literal"); + Assert.True(int.Parse(versionMatch.Groups[1].Value) >= 64); var start = initializer.IndexOf("if (fromVersion < 64)", StringComparison.Ordinal); Assert.True(start >= 0, "DuckDbInitializer has no v64 block"); var block = initializer[start..]; diff --git a/Darling/Darling.Tests/QueryStoreLivenessHotTouchLiveTests.cs b/Darling/Darling.Tests/QueryStoreLivenessHotTouchLiveTests.cs index f7218bc83..be7b65088 100644 --- a/Darling/Darling.Tests/QueryStoreLivenessHotTouchLiveTests.cs +++ b/Darling/Darling.Tests/QueryStoreLivenessHotTouchLiveTests.cs @@ -76,20 +76,22 @@ public void TheProbeCarriesThisRungsSentinel_AndTheArmSitsBelowTheCurrentTop() var arity = method.GetParameters().Length; Assert.Equal("hasHotLivenessTouch", method.GetParameters()[ProbeOrdinal].Name); - /* Every rung above this one (V150's hasCollectionLogWatermarkAndJobHistoryIndexes) must also be - false, or the map finds the newer arm first and this assertion is checking the wrong rung's - fallthrough. */ + /* Every rung above this one (V150's hasCollectionLogWatermarkAndJobHistoryIndexes, V151's + hasAgGroupId) must also be false, or the map finds a newer arm first and this assertion is + checking the wrong rung's fallthrough. */ var all = Enumerable.Repeat((object)true, arity).ToArray(); var behind = (object[])all.Clone(); - behind[ProbeOrdinal] = false; - behind[arity - 1] = false; + for (var i = ProbeOrdinal; i < arity; i++) + { + behind[i] = false; + } Assert.Equal(PreviousVersion, (int)method.Invoke(null, behind)!); - /* V150 (#4469, #4477) is now the top rung, so this arm no longer needs to be the LAST one — it only - has to sit below the current top's arm, which is what the ladder-dense invariant above already + /* V151 (#4475) is now the top rung, so this arm no longer needs to be the LAST one — it only has + to sit below the current top's arm, which is what the ladder-dense invariant above already guarantees is registered ahead of it. */ var thisArm = viewer.IndexOf("if (hasHotLivenessTouch)", StringComparison.Ordinal); - var topArm = viewer.IndexOf("if (hasCollectionLogWatermarkAndJobHistoryIndexes)", StringComparison.Ordinal); + var topArm = viewer.IndexOf("if (hasAgGroupId)", StringComparison.Ordinal); Assert.True(thisArm >= 0, "the viewer has no V149 sentinel arm — a fully-migrated store would map one rung short"); Assert.True(topArm >= 0 && topArm < thisArm, "the current top rung's arm must sit above the V149 arm"); Assert.Contains( diff --git a/Darling/Darling.Tests/ViewerAvailabilityGroupsTests.cs b/Darling/Darling.Tests/ViewerAvailabilityGroupsTests.cs index 8c068c071..d9a64e015 100644 --- a/Darling/Darling.Tests/ViewerAvailabilityGroupsTests.cs +++ b/Darling/Darling.Tests/ViewerAvailabilityGroupsTests.cs @@ -473,4 +473,96 @@ public void Counts_FortyTwoDistinctRdsAg0Instances_AreFortyTwoGroups() Assert.Equal(42, servers); Assert.Equal(42, views); } + + /* ────────────────── V151/#4475: the group_id overload's count rule ────────────────── */ + + [Fact] + public void CountDistinctGroups_TwoSecondariesOfOneAgSameGroupId_DisjointReplicaSets_IsOne() + { + /* The exact case a group_id closes: two monitored SECONDARIES of one real AG, its primary unmonitored. + Each reports only itself (sys.dm_hadr_availability_replica_states' local-only rule), so their replica + sets are DISJOINT and the name+overlap rule alone would count 2. Carrying the same group_id on both + unions them into 1. */ + var groups = AgTopology.CountDistinctGroups(new (string?, IEnumerable, string?)[] + { + ("RDSAG0", new string?[] { "S1" }, "GROUP-GUID-1"), + ("RDSAG0", new string?[] { "S2" }, "GROUP-GUID-1"), + }); + + Assert.Equal(1, groups); + } + + [Fact] + public void CountDistinctGroups_FortyTwoDistinctGroupIds_AreFortyTwoGroups() + { + /* 42 members, same AG name (RDSAG0), each carrying its OWN distinct group_id -- id-based union never + collapses them, same as the pre-existing disjoint-replica-set case, now proven on the id path too. */ + var members = new List<(string?, IEnumerable, string?)>(); + for (var i = 1; i <= 42; i++) + { + members.Add(("RDSAG0", new string?[] { $"NODE{i}A", $"NODE{i}B" }, $"GROUP-GUID-{i}")); + } + + var groups = AgTopology.CountDistinctGroups(members); + + Assert.Equal(42, groups); + } + + [Fact] + public void CountDistinctGroups_ExistingIdLessPins_AreUnchanged() + { + /* (c): the pre-existing id-less 2-tuple overload still behaves exactly as before -- calling it directly + (as every pre-V151 caller does) with no group_id anywhere reproduces the pre-existing #4475 41/42 + result unchanged. */ + var members = new List<(string?, IEnumerable)>(); + for (var i = 1; i <= 42; i++) + { + members.Add(("RDSAG0", new string?[] { $"NODE{i}A", $"NODE{i}B" })); + } + + Assert.Equal(42, AgTopology.CountDistinctGroups(members)); + } + + [Fact] + public void CountDistinctGroups_MixedWithIdAndIdLess_SameNameOverlappingReplicas_IsOne() + { + /* (d): a with-id member and an id-less member of the SAME name whose replicas overlap still union -- the + same AG, seen once before the V151 upgrade landed on that reporter (no group_id yet) and once after + (group_id now populated). */ + var groups = AgTopology.CountDistinctGroups(new (string?, IEnumerable, string?)[] + { + ("RDSAG0", new string?[] { "P", "S1" }, "GROUP-GUID-1"), + ("RDSAG0", new string?[] { "S1" }, null), + }); + + Assert.Equal(1, groups); + } + + [Fact] + public void CountDistinctGroups_TwoWithIdDifferentIds_SameNameOverlappingReplicas_IsTwo() + { + /* (e): two with-id members, DIFFERENT ids, same name, OVERLAPPING replicas -- a group_id is definitive, + so the name+overlap match that would otherwise union them is overridden. Never 1. */ + var groups = AgTopology.CountDistinctGroups(new (string?, IEnumerable, string?)[] + { + ("RDSAG0", new string?[] { "P", "S1" }, "GROUP-GUID-1"), + ("RDSAG0", new string?[] { "S1" }, "GROUP-GUID-2"), + }); + + Assert.Equal(2, groups); + } + + [Fact] + public void CountDistinctGroups_WithIdAndIdLess_SameNameDisjointReplicas_IsTwo() + { + /* (f): a with-id member and an id-less member share a name but their replicas are DISJOINT -- with no id + on one side to union on, and no replica overlap either, they stay 2 separate groups. */ + var groups = AgTopology.CountDistinctGroups(new (string?, IEnumerable, string?)[] + { + ("RDSAG0", new string?[] { "P", "S1" }, "GROUP-GUID-1"), + ("RDSAG0", new string?[] { "S2" }, null), + }); + + Assert.Equal(2, groups); + } } diff --git a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingAgReader.cs b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingAgReader.cs index 29fe53d69..a797872b8 100644 --- a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingAgReader.cs +++ b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingAgReader.cs @@ -179,6 +179,9 @@ the least-lagging tied groups first rather than an arbitrary alphabetical tail. ServerId = first.ServerId, ServerName = first.ServerName, AgName = first.AgName, + /* Every replica row of one AG carries the SAME group_id (V151); tolerate a row or two still + NULL mid-upgrade rather than requiring every row in the group to agree. */ + GroupId = replicaViews.Select(r => r.GroupId).FirstOrDefault(g => !string.IsNullOrWhiteSpace(g)), CollectionTime = first.CollectionTime, DatabaseCollectionTime = dbRows is { Count: > 0 } ? dbRows[0].CollectionTime : null, PrimaryReplica = replicaViews.FirstOrDefault(r => r.IsPrimary)?.ReplicaServerName, @@ -220,7 +223,7 @@ is measured before that tool's own limit cuts. A caller reading distinct_ag_coun returns local information only off the primary), so exact-set identity would double-count it. Shares AgTopology's counting helper so the viewer and this read cannot drift back apart. */ DistinctAgCount = AgTopology.CountDistinctGroups( - groups.Select(g => (g.AgName, (IEnumerable)g.Replicas.Select(r => r.ReplicaServerName)))), + groups.Select(g => (g.AgName, (IEnumerable)g.Replicas.Select(r => r.ReplicaServerName), g.GroupId))), WorstSeverity = groups.Count == 0 ? HealthSeverity.Unknown : groups.Max(g => g.Severity), AvailabilityGroups = pagedGroups, GroupsReturned = pagedGroups.Count, @@ -290,6 +293,7 @@ private static AgReplicaView ToReplicaView(ReplicaRow row) AvailabilityModeDesc = row.AvailabilityModeDesc, FailoverModeDesc = row.FailoverModeDesc, EndpointUrl = row.EndpointUrl, + GroupId = row.GroupId, }; } @@ -470,7 +474,8 @@ private static async Task> ReadReplicasAsync( rows.Add(new ReplicaRow( row.ServerId, row.ServerName, row.CollectionTime, row.AgName, row.ReplicaServerName, row.RoleDesc, row.IsLocal, row.OperationalStateDesc, row.ConnectedStateDesc, row.RecoveryHealthDesc, - row.SynchronizationHealthDesc, row.AvailabilityModeDesc, row.FailoverModeDesc, row.EndpointUrl)); + row.SynchronizationHealthDesc, row.AvailabilityModeDesc, row.FailoverModeDesc, row.EndpointUrl, + row.GroupId)); } return rows; @@ -488,7 +493,7 @@ private static async Task> ReadDatabasesAsync( row.ServerId, row.ServerName, row.CollectionTime, row.AgName, row.DatabaseName, row.ReplicaServerName, row.IsLocal, row.SynchronizationStateDesc, row.LastHardenedLsn, row.LastCommitLsn, row.LogSendQueueSize, row.RedoQueueSize, row.LogSendRate, row.RedoRate, row.IsSuspended, row.SuspendReasonDesc, - row.AvailabilityModeDesc, row.SecondaryLagSeconds)); + row.AvailabilityModeDesc, row.SecondaryLagSeconds, row.GroupId)); } return rows; @@ -511,7 +516,8 @@ internal readonly record struct ReplicaRow( string? SynchronizationHealthDesc, string? AvailabilityModeDesc, string? FailoverModeDesc, - string? EndpointUrl); + string? EndpointUrl, + string? GroupId); /// One database-grain row, exactly as collect.ag_database_replica_states stores it. Queue sizes /// are KB and rates KB/s (the DMV's units), both instantaneous gauges rather than counters. @@ -533,7 +539,8 @@ internal readonly record struct DatabaseRow( bool? IsSuspended, string? SuspendReasonDesc, string? AvailabilityModeDesc, - long? SecondaryLagSeconds); + long? SecondaryLagSeconds, + string? GroupId); } /// One replica inside a group, pre-banded. Every *_severity is derived server-side (R1). @@ -568,6 +575,10 @@ public sealed class AgReplicaView [JsonPropertyName("availability_mode")] public string? AvailabilityModeDesc { get; init; } [JsonPropertyName("failover_mode")] public string? FailoverModeDesc { get; init; } [JsonPropertyName("endpoint_url")] public string? EndpointUrl { get; init; } + + /// sys.availability_groups.group_id as text (V151, #4475). Null on a row collected before + /// this column existed. + [JsonPropertyName("group_id")] public string? GroupId { get; init; } } /// One database-on-a-replica row inside a group, pre-banded. @@ -615,6 +626,11 @@ public sealed class AvailabilityGroupView [JsonPropertyName("server_id")] public int ServerId { get; init; } [JsonPropertyName("ag_name")] public string? AgName { get; init; } + /// sys.availability_groups.group_id as text (V151, #4475) — the same GUID the engine stamps + /// on every replica of this AG, taken from whichever replica row carried it. Null when every replica row in + /// this group predates V151. + [JsonPropertyName("group_id")] public string? GroupId { get; init; } + /// When the reporting server's newest REPLICA-grain snapshot was taken (naive UTC). The collectors /// write nothing for a server with no AGs, so a group that stops refreshing keeps its last instant here — /// surfaced rather than silently aged out. diff --git a/Darling/PerformanceMonitor.Darling.Storage/DarlingAgStatesReader.cs b/Darling/PerformanceMonitor.Darling.Storage/DarlingAgStatesReader.cs index 6ea4db5ef..40c009ef2 100644 --- a/Darling/PerformanceMonitor.Darling.Storage/DarlingAgStatesReader.cs +++ b/Darling/PerformanceMonitor.Darling.Storage/DarlingAgStatesReader.cs @@ -143,7 +143,8 @@ LIMIT 1 r.synchronization_health_desc, r.availability_mode_desc, r.failover_mode_desc, - r.endpoint_url + r.endpoint_url, + r.group_id FROM ag_replica_states AS r JOIN ( @@ -186,7 +187,8 @@ AND s.is_enabled d.is_suspended, d.suspend_reason_desc, d.availability_mode_desc, - d.secondary_lag_seconds + d.secondary_lag_seconds, + d.group_id FROM ag_database_replica_states AS d JOIN ( @@ -247,7 +249,8 @@ public readonly record struct ReplicaStateRow( string? SynchronizationHealthDesc, string? AvailabilityModeDesc, string? FailoverModeDesc, - string? EndpointUrl); + string? EndpointUrl, + string? GroupId); /// One database-grain row, exactly as ag_database_replica_states stores it. Queue sizes /// are KB and rates KB/s (the DMV's units), both instantaneous gauges rather than counters. @@ -269,7 +272,8 @@ public readonly record struct DatabaseReplicaStateRow( bool? IsSuspended, string? SuspendReasonDesc, string? AvailabilityModeDesc, - long? SecondaryLagSeconds); + long? SecondaryLagSeconds, + string? GroupId); /* ─────────────────────────── reads ─────────────────────────── */ @@ -415,7 +419,8 @@ private static async Task> ReadReplicaStatesAsync( Text(reader, 10), Text(reader, 11), Text(reader, 12), - Text(reader, 13))); + Text(reader, 13), + Text(reader, 14))); } return rows; @@ -450,7 +455,8 @@ private static async Task> ReadDatabaseReplicaStat Flag(reader, 14), Text(reader, 15), Text(reader, 16), - Count(reader, 17))); + Count(reader, 17), + Text(reader, 18))); } return rows; diff --git a/Darling/PerformanceMonitor.Darling.Storage/PgMigrations.cs b/Darling/PerformanceMonitor.Darling.Storage/PgMigrations.cs index 15af21940..467d295df 100644 --- a/Darling/PerformanceMonitor.Darling.Storage/PgMigrations.cs +++ b/Darling/PerformanceMonitor.Darling.Storage/PgMigrations.cs @@ -230,6 +230,7 @@ holds the ordering. */ new Migration(148, "read-latency", V148Sql), new Migration(149, "query-store-liveness-hot-touch", V149Sql), new Migration(150, "collection-log-watermark-and-job-history-indexes", V150Sql), + new Migration(151, "ag-group-id", V151Sql), }; /// @@ -2403,6 +2404,48 @@ CREATE INDEX IF NOT EXISTS idx_collection_log_watermark CREATE INDEX IF NOT EXISTS idx_job_history_server_run ON collect.job_history (server_id, run_datetime DESC, instance_id DESC);"; + /// + /// V151 — sys.availability_groups.group_id on both AG collector tables (#4475): the GUID the engine + /// stamps identically on every replica of one Availability Group, appended LAST as text (the collector + /// column vocabulary has no uuid type; AgDatabaseReplicaStatesCollector's last_hardened_lsn / + /// last_commit_lsn already store a wide identifier the same way). + /// uses it to close the one gap the name-plus-replica-overlap rule (#4475) could not: two monitored + /// SECONDARIES of one AG, with its primary unmonitored, share no replica name with each other (each reports + /// only itself under sys.dm_hadr_availability_replica_states's local-only rule) and so counted as two + /// groups. The group_id is the same on both replicas' rows, so a member carrying one groups by it exactly; + /// a member from a row collected before this rung carries none and falls back to the pre-existing name + + /// overlap rule among the other id-less members; and a with-id member and a without-id member of the same + /// name whose replica sets overlap still join (the same AG, seen before and after the upgrade landed on that + /// reporter). Stated in 's own doc, not restated as a second + /// source of truth here. + /// + /// Nullable, no DEFAULT, no backfill, the V127/V128/V132/V133/V150 shape for every column-adding + /// rung on a collector table: a row collected before this rung never asked the engine for its AG's + /// group_id, and NULL is the honest value — a reader treats it as "fall back to the pre-#4475 name + overlap + /// rule", exactly today's behavior. Both tables are compressed hypertables on the fleet; a nullable, + /// default-less ADD COLUMN is catalog-only in PostgreSQL and TimescaleDB accepts it on a compressed + /// hypertable with a compression policy attached, the shape V127/V128/V132/V133 used and verified live each + /// time. No view to refresh: the AG collector tables have been view-less since V34 (no v_ passthrough + /// was ever created for either), so this rung is two ALTERs and nothing else. + /// + /// A fresh store gets the column from the generated CREATE TABLE + /// ( / carry it, + /// appended last in PayloadColumns) and the ALTER no-ops there — the V101 rule, pinned by + /// PgSchemaGeneratorTests reconstructing the current shape from V34/V36/V37/this rung and comparing it + /// to the generator's current output. + /// + /// Lite's DuckDB twin gets the same column the same way (schema version bump, additive ALTER TABLE + /// ... ADD COLUMN IF NOT EXISTS), and the shared is what + /// both Lite's AG tab and Darling's Viewer/MCP/web reads call, so the fallback rule cannot drift apart + /// between stores. + /// + private const string V151Sql = @" +ALTER TABLE collect.ag_replica_states + ADD COLUMN IF NOT EXISTS group_id text; + +ALTER TABLE collect.ag_database_replica_states + ADD COLUMN IF NOT EXISTS group_id text;"; + /// /// V2 — the service's observability store: the servers registry (upserted on every /// successful connect) and the per-run collection_log. Column names deliberately mirror diff --git a/Darling/PerformanceMonitor.Darling.Storage/StorageVersion.cs b/Darling/PerformanceMonitor.Darling.Storage/StorageVersion.cs index 6dad0263d..ffab2543d 100644 --- a/Darling/PerformanceMonitor.Darling.Storage/StorageVersion.cs +++ b/Darling/PerformanceMonitor.Darling.Storage/StorageVersion.cs @@ -16,5 +16,5 @@ namespace PerformanceMonitor.Darling.Storage; /// public static class StorageVersion { - public const int SchemaVersion = 150; + public const int SchemaVersion = 151; } diff --git a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.AvailabilityGroups.cs b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.AvailabilityGroups.cs index 0b68a0ff9..2d6c7b0cf 100644 --- a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.AvailabilityGroups.cs +++ b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.AvailabilityGroups.cs @@ -72,6 +72,7 @@ private async Task> ReadAgReplicasAsync(DateTime nowU AvailabilityModeDesc = row.AvailabilityModeDesc, FailoverModeDesc = row.FailoverModeDesc, EndpointUrl = row.EndpointUrl, + GroupId = row.GroupId, }); } @@ -105,6 +106,7 @@ private async Task> ReadAgDatabasesAsync(DateTime no SuspendReasonDesc = row.SuspendReasonDesc, AvailabilityModeDesc = row.AvailabilityModeDesc, SecondaryLagSeconds = row.SecondaryLagSeconds, + GroupId = row.GroupId, }); } diff --git a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.cs b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.cs index 3847ca489..2c22265a1 100644 --- a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.cs +++ b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.cs @@ -989,7 +989,12 @@ AND NOT EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexn read by any viewer surface, so this gate rests on the standing invariant alone. Named only in this probe line, never in prose, per the V71 finding. */ (EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexname = 'idx_collection_log_watermark') - AND EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexname = 'idx_job_history_server_run'))"; + AND EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexname = 'idx_job_history_server_run')), + /* V151 (#4475) probes a COLUMN for the V36/V37 reason: ag_replica_states has existed since V34, so + table existence cannot separate the rungs. It is not yet read by any viewer surface, so this gate + rests on the standing invariant alone. Named only in this probe line, never in prose, per the V71 + finding. */ + EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name = 'ag_replica_states' AND column_name = 'group_id')"; /// The store schema version this viewer build requires — the highest migration it knows /// (). The connect-time gate blocks a store below this. @@ -1011,7 +1016,7 @@ AND NOT EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexn await using var reader = await command.ExecuteReaderAsync(cancellationToken); if (await reader.ReadAsync(cancellationToken)) { - return MapProbedSchemaVersion(reader.GetBoolean(0), reader.GetBoolean(1), reader.GetBoolean(2), reader.GetBoolean(3), reader.GetBoolean(4), reader.GetBoolean(5), reader.GetBoolean(6), reader.GetBoolean(7), reader.GetBoolean(8), reader.GetBoolean(9), reader.GetBoolean(10), reader.GetBoolean(11), reader.GetBoolean(12), reader.GetBoolean(13), reader.GetBoolean(14), reader.GetBoolean(15), reader.GetBoolean(16), reader.GetBoolean(17), reader.GetBoolean(18), reader.GetBoolean(19), reader.GetBoolean(20), reader.GetBoolean(21), reader.GetBoolean(22), reader.GetBoolean(23), reader.GetBoolean(24), reader.GetBoolean(25), reader.GetBoolean(26), reader.GetBoolean(27), reader.GetBoolean(28), reader.GetBoolean(29), reader.GetBoolean(30), reader.GetBoolean(31), reader.GetBoolean(32), reader.GetBoolean(33), reader.GetBoolean(34), reader.GetBoolean(35), reader.GetBoolean(36), reader.GetBoolean(37), reader.GetBoolean(38), reader.GetBoolean(39), reader.GetBoolean(40), reader.GetBoolean(41), reader.GetBoolean(42), reader.GetBoolean(43), reader.GetBoolean(44), reader.GetBoolean(45), reader.GetBoolean(46), reader.GetBoolean(47), reader.GetBoolean(48), reader.GetBoolean(49), reader.GetBoolean(50), reader.GetBoolean(51), reader.GetBoolean(52), reader.GetBoolean(53), reader.GetBoolean(54), reader.GetBoolean(55), reader.GetBoolean(56), reader.GetBoolean(57), reader.GetBoolean(58), reader.GetBoolean(59), reader.GetBoolean(60), reader.GetBoolean(61), reader.GetBoolean(62), reader.GetBoolean(63), reader.GetBoolean(64), reader.GetBoolean(65), reader.GetBoolean(66), reader.GetBoolean(67), reader.GetBoolean(68), reader.GetBoolean(69), reader.GetBoolean(70), reader.GetBoolean(71), reader.GetBoolean(72), reader.GetBoolean(73), reader.GetBoolean(74), reader.GetBoolean(75), reader.GetBoolean(76), reader.GetBoolean(77), reader.GetBoolean(78), reader.GetBoolean(79), reader.GetBoolean(80), reader.GetBoolean(81), reader.GetBoolean(82), reader.GetBoolean(83), reader.GetBoolean(84), reader.GetBoolean(85), reader.GetBoolean(86), reader.GetBoolean(87), reader.GetBoolean(88), reader.GetBoolean(89), reader.GetBoolean(90), reader.GetBoolean(91), reader.GetBoolean(92), reader.GetBoolean(93), reader.GetBoolean(94), reader.GetBoolean(95), reader.GetBoolean(96), reader.GetBoolean(97), reader.GetBoolean(98), reader.GetBoolean(99), reader.GetBoolean(100), reader.GetBoolean(101), reader.GetBoolean(102), reader.GetBoolean(103), reader.GetBoolean(104), reader.GetBoolean(105), reader.GetBoolean(106), reader.GetBoolean(107), reader.GetBoolean(108), reader.GetBoolean(109), reader.GetBoolean(110), reader.GetBoolean(111), reader.GetBoolean(112), reader.GetBoolean(113), reader.GetBoolean(114), reader.GetBoolean(115), reader.GetBoolean(116), reader.GetBoolean(117), reader.GetBoolean(118), reader.GetBoolean(119), reader.GetBoolean(120), reader.GetBoolean(121), reader.GetBoolean(122), reader.GetBoolean(123), reader.GetBoolean(124), reader.GetBoolean(125)); + return MapProbedSchemaVersion(reader.GetBoolean(0), reader.GetBoolean(1), reader.GetBoolean(2), reader.GetBoolean(3), reader.GetBoolean(4), reader.GetBoolean(5), reader.GetBoolean(6), reader.GetBoolean(7), reader.GetBoolean(8), reader.GetBoolean(9), reader.GetBoolean(10), reader.GetBoolean(11), reader.GetBoolean(12), reader.GetBoolean(13), reader.GetBoolean(14), reader.GetBoolean(15), reader.GetBoolean(16), reader.GetBoolean(17), reader.GetBoolean(18), reader.GetBoolean(19), reader.GetBoolean(20), reader.GetBoolean(21), reader.GetBoolean(22), reader.GetBoolean(23), reader.GetBoolean(24), reader.GetBoolean(25), reader.GetBoolean(26), reader.GetBoolean(27), reader.GetBoolean(28), reader.GetBoolean(29), reader.GetBoolean(30), reader.GetBoolean(31), reader.GetBoolean(32), reader.GetBoolean(33), reader.GetBoolean(34), reader.GetBoolean(35), reader.GetBoolean(36), reader.GetBoolean(37), reader.GetBoolean(38), reader.GetBoolean(39), reader.GetBoolean(40), reader.GetBoolean(41), reader.GetBoolean(42), reader.GetBoolean(43), reader.GetBoolean(44), reader.GetBoolean(45), reader.GetBoolean(46), reader.GetBoolean(47), reader.GetBoolean(48), reader.GetBoolean(49), reader.GetBoolean(50), reader.GetBoolean(51), reader.GetBoolean(52), reader.GetBoolean(53), reader.GetBoolean(54), reader.GetBoolean(55), reader.GetBoolean(56), reader.GetBoolean(57), reader.GetBoolean(58), reader.GetBoolean(59), reader.GetBoolean(60), reader.GetBoolean(61), reader.GetBoolean(62), reader.GetBoolean(63), reader.GetBoolean(64), reader.GetBoolean(65), reader.GetBoolean(66), reader.GetBoolean(67), reader.GetBoolean(68), reader.GetBoolean(69), reader.GetBoolean(70), reader.GetBoolean(71), reader.GetBoolean(72), reader.GetBoolean(73), reader.GetBoolean(74), reader.GetBoolean(75), reader.GetBoolean(76), reader.GetBoolean(77), reader.GetBoolean(78), reader.GetBoolean(79), reader.GetBoolean(80), reader.GetBoolean(81), reader.GetBoolean(82), reader.GetBoolean(83), reader.GetBoolean(84), reader.GetBoolean(85), reader.GetBoolean(86), reader.GetBoolean(87), reader.GetBoolean(88), reader.GetBoolean(89), reader.GetBoolean(90), reader.GetBoolean(91), reader.GetBoolean(92), reader.GetBoolean(93), reader.GetBoolean(94), reader.GetBoolean(95), reader.GetBoolean(96), reader.GetBoolean(97), reader.GetBoolean(98), reader.GetBoolean(99), reader.GetBoolean(100), reader.GetBoolean(101), reader.GetBoolean(102), reader.GetBoolean(103), reader.GetBoolean(104), reader.GetBoolean(105), reader.GetBoolean(106), reader.GetBoolean(107), reader.GetBoolean(108), reader.GetBoolean(109), reader.GetBoolean(110), reader.GetBoolean(111), reader.GetBoolean(112), reader.GetBoolean(113), reader.GetBoolean(114), reader.GetBoolean(115), reader.GetBoolean(116), reader.GetBoolean(117), reader.GetBoolean(118), reader.GetBoolean(119), reader.GetBoolean(120), reader.GetBoolean(121), reader.GetBoolean(122), reader.GetBoolean(123), reader.GetBoolean(124), reader.GetBoolean(125), reader.GetBoolean(126)); } return null; @@ -1036,7 +1041,7 @@ AND NOT EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexn /// is unit-tested without a live store; any schema bump past the newest arm trips the pinning test that keeps /// this in step with . /// - internal static int MapProbedSchemaVersion(bool hasConfigControlPlane, bool hasAlertDeliveryOverride, bool hasAnalysisState, bool hasAlertTuningKnobs, bool hasDefaultTraceEvents, bool hasIndexObjectStatsLatestIndex, bool hasCollectionLogHypertableOrPlainPg, bool hasJobHistory, bool hasAgentStatus, bool hasGenericWebhook, bool hasDeadlocksDatabaseName, bool hasQueryStoreReplicaRole, bool hasLongQueryCompletions, bool hasWebDashboardConfig, bool hasCustomViews, bool hasServerTags, bool hasConnectionRefireKnobs = false, bool hasAgCollectors = false, bool hasAgAlertKnobs = false, bool hasAgLatencyColumns = false, bool hasAgDisconnectRefire = false, bool hasPayloadDimensions = false, bool hasDimFloorIndexes = false, bool hasBlockingWaitThreshold = false, bool hasQueryStoreIntervalIdentity = false, bool hasPagerDutyWebhook = false, bool hasPagerDutyProxy = false, bool hasCollectorState = false, bool hasPlanCorrection = false, bool hasPvsStats = false, bool hasPvsPressureKnobs = false, bool hasDatabaseStateAlert = false, bool hasServerTagColour = false, bool hasQueryStatsHostObject = false, bool hasFindingDrillDown = false, bool hasStoreMetrics = false, bool hasPlanDimGzip = false, bool hasSelfAlertKnobs = false, bool hasJobMetricsColumns = false, bool hasJobCadenceKnob = false, bool hasBackfillSwitch = false, bool hasCollectorMemoryKnobs = false, bool hasDatabaseStateEdgeMemory = false, bool hasIncidentOccurrences = false, bool hasPlanXmlCompressionKnob = false, bool hasMonitoredServerEngine = false, bool hasPgBlockingEdges = false, bool hasQueryStorePlanMap = false, bool hasPgStatementText = false, bool hasQueryStoreText = false, bool hasPlanContentRetentionKnob = false, bool hasQueryStoreHealth = false, bool hasQueryStoreTextHash = false, bool hasComposeTimeoutKnob = false, bool hasFileGrowthAlert = false, bool hasCollectionLogFanoutRollup = false, bool hasTempDbMaxSize = false, bool hasServerEngineKind = false, bool hasPgDatabaseStats = false, bool hasPgIndexUsageStats = false, bool hasPgTableBloatStats = false, bool hasPgSessionStates = false, bool hasPgPlanCaptureReadiness = false, bool hasPgWriteStats = false, bool hasPgExtensionAvailability = false, bool hasPgLockStats = false, bool hasPgColumnStats = false, bool hasPgReplicationStats = false, bool hasPgBufferUsage = false, bool hasPgIndexBloat = false, bool hasPgPerDatabaseAttribution = false, bool hasPgWaitSampling = false, bool hasPgKernelStats = false, bool hasPgPredicateStats = false, bool hasPgPlanCapture = false, bool hasPgMajorVersion = false, bool hasPg18IoBytes = false, bool hasPgServerConfig = false, bool hasPgDeadlocks = false, bool hasPgDeadlockIdentity = false, bool hasCollectorCost = false, bool hasPgCpuUtilization = false, bool hasPlanForceActions = false, bool hasCollectionLogPhaseSplit = false, bool hasCollectionLogDrainForensics = false, bool hasCollectionLogFetchPhaseSums = false, bool hasStoreLogSelfMonitoring = false, bool hasCollectorStallProbes = false, bool hasRemediationCredentialAndActor = false, bool hasPgIndexBloatEstimate = false, bool hasPgCpuCapacityHeadroom = false, bool hasCustomAlertCore = false, bool hasMuteRuleReloadBeacon = false, bool hasBuiltinAlertPersistence = false, bool hasRetentionHoldRatioKnobs = false, bool hasDeadlockRateBandKnobs = false, bool hasOversizedPlanBacklog = false, bool hasPgAlertCountKnobs = false, bool hasFleetSweepState = false, bool hasFleetSweepCadenceKnobs = false, bool hasCollectorScheduleDatabases = false, bool hasSelfDiskWarnGbFloor = false, bool hasDeltaFamilyIntervalColumns = false, bool hasDeltaFamilyIntervalCompletion = false, bool hasPgLogEvents = false, bool hasPgLogEventMetrics = false, bool hasNotificationRoutes = false, bool hasPerfmonCounterType = false, bool hasPgNumbackendsAndSampledMs = false, bool hasTimeHonesty = false, bool hasLrqExclusionKnob = false, bool hasPgDatabaseSizeStatsAndHostMemory = false, bool hasQsCaptureModeRouteKnobToast = false, bool hasPgServerConfigDatabaseRoleOverrides = false, bool hasPostmasterStartTime = false, bool hasCheckpointsTimed = false, bool hasCollectionCaveats = false, bool hasIndexObjectStatsServerTimeIndex = false, bool hasQueryStoreIntervalLatest = false, bool hasRawChunkIntervalRungHistory = false, bool hasQueryStoreIntervalWide = false, bool hasManagedConfVerdicts = false, bool hasComposeTimeoutSixty = false, bool hasReadLatency = false, bool hasHotLivenessTouch = false, bool hasCollectionLogWatermarkAndJobHistoryIndexes = false) + internal static int MapProbedSchemaVersion(bool hasConfigControlPlane, bool hasAlertDeliveryOverride, bool hasAnalysisState, bool hasAlertTuningKnobs, bool hasDefaultTraceEvents, bool hasIndexObjectStatsLatestIndex, bool hasCollectionLogHypertableOrPlainPg, bool hasJobHistory, bool hasAgentStatus, bool hasGenericWebhook, bool hasDeadlocksDatabaseName, bool hasQueryStoreReplicaRole, bool hasLongQueryCompletions, bool hasWebDashboardConfig, bool hasCustomViews, bool hasServerTags, bool hasConnectionRefireKnobs = false, bool hasAgCollectors = false, bool hasAgAlertKnobs = false, bool hasAgLatencyColumns = false, bool hasAgDisconnectRefire = false, bool hasPayloadDimensions = false, bool hasDimFloorIndexes = false, bool hasBlockingWaitThreshold = false, bool hasQueryStoreIntervalIdentity = false, bool hasPagerDutyWebhook = false, bool hasPagerDutyProxy = false, bool hasCollectorState = false, bool hasPlanCorrection = false, bool hasPvsStats = false, bool hasPvsPressureKnobs = false, bool hasDatabaseStateAlert = false, bool hasServerTagColour = false, bool hasQueryStatsHostObject = false, bool hasFindingDrillDown = false, bool hasStoreMetrics = false, bool hasPlanDimGzip = false, bool hasSelfAlertKnobs = false, bool hasJobMetricsColumns = false, bool hasJobCadenceKnob = false, bool hasBackfillSwitch = false, bool hasCollectorMemoryKnobs = false, bool hasDatabaseStateEdgeMemory = false, bool hasIncidentOccurrences = false, bool hasPlanXmlCompressionKnob = false, bool hasMonitoredServerEngine = false, bool hasPgBlockingEdges = false, bool hasQueryStorePlanMap = false, bool hasPgStatementText = false, bool hasQueryStoreText = false, bool hasPlanContentRetentionKnob = false, bool hasQueryStoreHealth = false, bool hasQueryStoreTextHash = false, bool hasComposeTimeoutKnob = false, bool hasFileGrowthAlert = false, bool hasCollectionLogFanoutRollup = false, bool hasTempDbMaxSize = false, bool hasServerEngineKind = false, bool hasPgDatabaseStats = false, bool hasPgIndexUsageStats = false, bool hasPgTableBloatStats = false, bool hasPgSessionStates = false, bool hasPgPlanCaptureReadiness = false, bool hasPgWriteStats = false, bool hasPgExtensionAvailability = false, bool hasPgLockStats = false, bool hasPgColumnStats = false, bool hasPgReplicationStats = false, bool hasPgBufferUsage = false, bool hasPgIndexBloat = false, bool hasPgPerDatabaseAttribution = false, bool hasPgWaitSampling = false, bool hasPgKernelStats = false, bool hasPgPredicateStats = false, bool hasPgPlanCapture = false, bool hasPgMajorVersion = false, bool hasPg18IoBytes = false, bool hasPgServerConfig = false, bool hasPgDeadlocks = false, bool hasPgDeadlockIdentity = false, bool hasCollectorCost = false, bool hasPgCpuUtilization = false, bool hasPlanForceActions = false, bool hasCollectionLogPhaseSplit = false, bool hasCollectionLogDrainForensics = false, bool hasCollectionLogFetchPhaseSums = false, bool hasStoreLogSelfMonitoring = false, bool hasCollectorStallProbes = false, bool hasRemediationCredentialAndActor = false, bool hasPgIndexBloatEstimate = false, bool hasPgCpuCapacityHeadroom = false, bool hasCustomAlertCore = false, bool hasMuteRuleReloadBeacon = false, bool hasBuiltinAlertPersistence = false, bool hasRetentionHoldRatioKnobs = false, bool hasDeadlockRateBandKnobs = false, bool hasOversizedPlanBacklog = false, bool hasPgAlertCountKnobs = false, bool hasFleetSweepState = false, bool hasFleetSweepCadenceKnobs = false, bool hasCollectorScheduleDatabases = false, bool hasSelfDiskWarnGbFloor = false, bool hasDeltaFamilyIntervalColumns = false, bool hasDeltaFamilyIntervalCompletion = false, bool hasPgLogEvents = false, bool hasPgLogEventMetrics = false, bool hasNotificationRoutes = false, bool hasPerfmonCounterType = false, bool hasPgNumbackendsAndSampledMs = false, bool hasTimeHonesty = false, bool hasLrqExclusionKnob = false, bool hasPgDatabaseSizeStatsAndHostMemory = false, bool hasQsCaptureModeRouteKnobToast = false, bool hasPgServerConfigDatabaseRoleOverrides = false, bool hasPostmasterStartTime = false, bool hasCheckpointsTimed = false, bool hasCollectionCaveats = false, bool hasIndexObjectStatsServerTimeIndex = false, bool hasQueryStoreIntervalLatest = false, bool hasRawChunkIntervalRungHistory = false, bool hasQueryStoreIntervalWide = false, bool hasManagedConfVerdicts = false, bool hasComposeTimeoutSixty = false, bool hasReadLatency = false, bool hasHotLivenessTouch = false, bool hasCollectionLogWatermarkAndJobHistoryIndexes = false, bool hasAgGroupId = false) { /* V71 (the PostgreSQL blocking-edges rung): a table-existence sentinel and now the newest-first arm. A collector table would ordinarily get no arm at all — see the V63-V69 note below — but the TOP @@ -1216,11 +1221,24 @@ the rung below and showing a spurious upgrade banner on a store that is current. The WPF viewer runs no analysis, so no viewer read names the new table; this arm exists so the version banner stays truthful, which is the only effect the rung has on the viewer. Named only in the probe line, not this prose, per the V71 finding. */ - /* V150 (#4469, #4477): two supporting indexes, idx_collection_log_watermark and - idx_job_history_server_run, and now the TOP rung, so a fully-migrated store maps to EXACTLY + /* V151 (#4475): ag_replica_states.group_id, a column that identifies which physical AG a replica + row belongs to, and now the TOP rung, so a fully-migrated store maps to EXACTLY StorageVersion.SchemaVersion rather than falling through to the rung below and showing a spurious upgrade banner on a store that is current. + The WPF viewer runs no analysis, so no viewer read names the new column; this arm exists so the + version banner stays truthful, which is the only effect the rung has on the viewer. Named only + in the probe line, not this prose, per the V71 finding. */ + if (hasAgGroupId) + { + return 151; + } + + /* V150 (#4469, #4477): two supporting indexes, idx_collection_log_watermark and + idx_job_history_server_run. Formerly the TOP rung — RequiredStoreSchemaVersion is + StorageVersion.SchemaVersion and a store below this arm now falls through to V149 instead of + stopping here. + The WPF viewer runs no analysis, so no viewer read names either index; this arm exists so the version banner stays truthful, which is the only effect the rung has on the viewer. Named only in the probe line, not this prose, per the V71 finding. */ diff --git a/Lite.Tests/AgCollectorDefinitionTests.cs b/Lite.Tests/AgCollectorDefinitionTests.cs index deb7636ae..d7cf7c5ee 100644 --- a/Lite.Tests/AgCollectorDefinitionTests.cs +++ b/Lite.Tests/AgCollectorDefinitionTests.cs @@ -48,7 +48,8 @@ public sealed class AgCollectorDefinitionTests availability_mode_desc = ar.availability_mode_desc, failover_mode_desc = ar.failover_mode_desc, endpoint_url = ar.endpoint_url, - is_local = ars.is_local + is_local = ars.is_local, + group_id = CONVERT(nvarchar(36), ag.group_id) FROM sys.availability_replicas AS ar JOIN sys.availability_groups AS ag ON ar.group_id = ag.group_id @@ -84,7 +85,8 @@ ORDER BY last_redone_time = hdrs.last_redone_time, last_received_time = hdrs.last_received_time, est_redo_completion_time_min = CONVERT(float, (hdrs.redo_queue_size * 1.0 / NULLIF(hdrs.redo_rate, 0)) / 60.0), - est_send_drain_time_min = CONVERT(float, (hdrs.log_send_queue_size * 1.0 / NULLIF(hdrs.log_send_rate, 0)) / 60.0) + est_send_drain_time_min = CONVERT(float, (hdrs.log_send_queue_size * 1.0 / NULLIF(hdrs.log_send_rate, 0)) / 60.0), + group_id = CONVERT(nvarchar(36), ag.group_id) FROM sys.dm_hadr_database_replica_states AS hdrs JOIN sys.availability_replicas AS ar ON hdrs.replica_id = ar.replica_id @@ -186,6 +188,9 @@ public void ReplicaPayloadColumns_AreInAppendOrder() /* #1696: appended LAST, so an upgraded store's ALTER lands it in the same physical position a fresh generated table puts it - which is what keeps the two provenances comparable. */ "is_local", + /* #4475: appended AFTER is_local, for the same reason - it lands in the same physical + position a fresh generated table puts it. */ + "group_id", }, names); @@ -227,9 +232,13 @@ public void DatabasePayloadColumns_AreInAppendOrder_WithTheDeclaredTypes() "last_received_time", "est_redo_completion_time_min", "est_send_drain_time_min", + /* #4475: appended LAST, same rationale as the replica table's group_id. */ + "group_id", }, columns.Select(c => c.Name).ToArray()); + Assert.Equal(CollectorColumnType.Varchar, columns.Single(c => c.Name == "group_id").Type); + /* The LSNs are numeric(25,0) in the DMV — far wider than BIGINT — so they are converted server-side and stored as text. Storing them numerically would silently overflow. */ Assert.Equal(CollectorColumnType.Varchar, columns.Single(c => c.Name == "last_hardened_lsn").Type); @@ -276,8 +285,8 @@ whole collection cycle rather than one column. */ public async Task ReplicaReadAsync_MapsColumns() { using var reader = new FakeCollectorDataReader( - new object[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE1.corp:5022", true }, - new object[] { "AG1", "NODE2", "SECONDARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE2.corp:5022", false }); + new object[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE1.corp:5022", true, "11111111-1111-1111-1111-111111111111" }, + new object[] { "AG1", "NODE2", "SECONDARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE2.corp:5022", false, "11111111-1111-1111-1111-111111111111" }); var context = CollectorTestContext.Make(new RecordingCollectorDeltaCalculator()); @@ -285,7 +294,7 @@ public async Task ReplicaReadAsync_MapsColumns() Assert.Equal(2, rows.Count); Assert.Equal( - new AgReplicaStatesCollector.Row("AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE1.corp:5022", true), + new AgReplicaStatesCollector.Row("AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", "TCP://NODE1.corp:5022", true, "11111111-1111-1111-1111-111111111111"), rows[0]); Assert.Equal("NODE2", rows[1].ReplicaServerName); Assert.Equal("SECONDARY", rows[1].RoleDesc); @@ -293,6 +302,8 @@ public async Task ReplicaReadAsync_MapsColumns() node's view of an AG that every node can see. */ Assert.True(rows[0].IsLocal); Assert.False(rows[1].IsLocal); + /* #4475: group_id maps like any other trailing text column. */ + Assert.Equal("11111111-1111-1111-1111-111111111111", rows[0].GroupId); } [Fact] @@ -302,8 +313,8 @@ public async Task ReplicaReadAsync_ToleratesNulls() state sys.availability_replicas serves only locally cached metadata — so every column can be null. A quorum-loss read must produce a row, not an InvalidCastException. */ using var reader = new FakeCollectorDataReader( - new object[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", DBNull.Value, DBNull.Value }, - new object[] { DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value }); + new object[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", DBNull.Value, DBNull.Value, DBNull.Value }, + new object[] { DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value }); var context = CollectorTestContext.Make(new RecordingCollectorDeltaCalculator()); @@ -317,6 +328,8 @@ state sys.availability_replicas serves only locally cached metadata — so every Assert.Null(rows[1].IsLocal); Assert.Equal("AG1", rows[0].AgName); Assert.Equal(default(AgReplicaStatesCollector.Row), rows[1]); + /* #4475: group_id is nullable, and NULL here (a pre-v65 row would carry none). */ + Assert.Null(rows[0].GroupId); } [Fact] @@ -335,6 +348,8 @@ public async Task DatabaseReadAsync_MapsColumns() they need not agree with the queue/rate columns above. Deliberately different values so a swapped mapping between the two fails rather than passing on a coincidence. */ 8.0d, 0.5d, + /* #4475: group_id, trailing. */ + "22222222-2222-2222-2222-222222222222", }); var context = CollectorTestContext.Make(new RecordingCollectorDeltaCalculator()); @@ -364,6 +379,8 @@ swapped mapping between the two fails rather than passing on a coincidence. */ Assert.Equal(8.0d, row.EstRedoCompletionTimeMin); Assert.Equal(0.5d, row.EstSendDrainTimeMin); + /* #4475: group_id maps like any other trailing text column. */ + Assert.Equal("22222222-2222-2222-2222-222222222222", row.GroupId); } [Fact] @@ -383,6 +400,8 @@ are null on a replica that is not reporting. The two drain estimates are ALWAYS DBNull.Value, DBNull.Value, "ASYNCHRONOUS_COMMIT", DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, DBNull.Value, + /* #4475: group_id, trailing, nullable. */ + DBNull.Value, }); var context = CollectorTestContext.Make(new RecordingCollectorDeltaCalculator()); @@ -408,6 +427,8 @@ are null on a replica that is not reporting. The two drain estimates are ALWAYS /* NULL, not 0 — an un-drainable queue must never read as "drains instantly". */ Assert.Null(row.EstRedoCompletionTimeMin); Assert.Null(row.EstSendDrainTimeMin); + /* #4475: group_id is nullable, and NULL here (a pre-v65 row would carry none). */ + Assert.Null(row.GroupId); } [Fact] @@ -433,12 +454,12 @@ public void WritePayload_EmitsPayloadOrder_AndTakesNoDeltas() var replicaWriter = new RecordingCollectorRowWriter(); AgReplicaStatesCollector.Instance.WritePayload( - new AgReplicaStatesCollector.Row("AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", null, true), + new AgReplicaStatesCollector.Row("AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", null, true, null), replicaWriter, context); Assert.Equal( - new object?[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", null, true }, + new object?[] { "AG1", "NODE1", "PRIMARY", "ONLINE", "CONNECTED", "ONLINE", "HEALTHY", "SYNCHRONOUS_COMMIT", "AUTOMATIC", null, true, null }, replicaWriter.Values); /* Modeled on a real measured sample from the Docker AG fixture: a SUSPEND_FROM_USER replica 62 s @@ -451,7 +472,7 @@ public void WritePayload_EmitsPayloadOrder_AndTakesNoDeltas() "AG1", "Orders", "NODE2", false, "SYNCHRONIZING", "1", "2", 4096L, 2048L, 512L, 256L, true, "SUSPEND_FROM_USER", "SYNCHRONOUS_COMMIT", 62L, new DateTime(2026, 7, 26, 12, 0, 0), new DateTime(2026, 7, 26, 12, 0, 1), new DateTime(2026, 7, 26, 11, 59, 55), new DateTime(2026, 7, 26, 12, 0, 2), - 8.0d, null), + 8.0d, null, null), databaseWriter, context); @@ -461,7 +482,7 @@ public void WritePayload_EmitsPayloadOrder_AndTakesNoDeltas() "AG1", "Orders", "NODE2", false, "SYNCHRONIZING", "1", "2", 4096L, 2048L, 512L, 256L, true, "SUSPEND_FROM_USER", "SYNCHRONOUS_COMMIT", 62L, new DateTime(2026, 7, 26, 12, 0, 0), new DateTime(2026, 7, 26, 12, 0, 1), new DateTime(2026, 7, 26, 11, 59, 55), new DateTime(2026, 7, 26, 12, 0, 2), - 8.0d, null, + 8.0d, null, null, }, databaseWriter.Values); diff --git a/Lite.Tests/DuckDbSchemaEquivalenceTests.cs b/Lite.Tests/DuckDbSchemaEquivalenceTests.cs index 50ac2298a..62d2cc1ce 100644 --- a/Lite.Tests/DuckDbSchemaEquivalenceTests.cs +++ b/Lite.Tests/DuckDbSchemaEquivalenceTests.cs @@ -211,6 +211,13 @@ fail until the oracle is extended by hand. That is the whole point of an oracle. /// DMV's type), nullable, trailing; NULL on every pre-v62 row is "type never recorded", which the readers /// classify by the #3702 name proxy as they did before the rung. 's v62 /// migration adds it to existing databases. + /// + /// #4475 / schema v65: group_id on ag_replica_states and + /// ag_database_replica_states — the Availability Group's engine-assigned GUID from + /// sys.availability_groups.group_id, the same value on every replica, so the distinct-group + /// count can tell same-named AGs apart. VARCHAR, nullable, trailing; NULL on every pre-v65 row, which + /// the readers fall back to the name-plus-replica-overlap heuristic for. 's + /// v65 migration adds it to existing databases. /// private static readonly HashSet IntentionalAppendedColumns = new(StringComparer.Ordinal) { @@ -227,6 +234,9 @@ fail until the oracle is extended by hand. That is the whole point of an oracle. "query_stats.statement_end_offset", /* v62 (#3653 A7): the counter's DMV type, so gauges stop being differenced. Appended, nullable INTEGER. */ "perfmon_stats.cntr_type", + /* v65 (#4475): the AG's group_id from sys.availability_groups, so same-named AGs can be told apart. Appended, nullable VARCHAR. */ + "ag_replica_states.group_id", + "ag_database_replica_states.group_id", }; [Fact] diff --git a/Lite.Tests/QsCaptureModeRungTests.cs b/Lite.Tests/QsCaptureModeRungTests.cs index ef6478f46..8fe66123d 100644 --- a/Lite.Tests/QsCaptureModeRungTests.cs +++ b/Lite.Tests/QsCaptureModeRungTests.cs @@ -162,7 +162,9 @@ state an existing database is in when the ALTER runs — the ADD COLUMN must suc using (var conn = new DuckDBConnection($"Data Source={dbPath}")) { await conn.OpenAsync(); - Assert.Equal(64L, Convert.ToInt64(await ScalarAsync(conn, "SELECT MAX(version) FROM schema_version"))); + /* Stays-true shape (#4475 raised CurrentSchemaVersion past 64): the climb reaches AT LEAST v64, + and lands exactly on CurrentSchemaVersion, whatever that is today. */ + Assert.True(DuckDbInitializer.CurrentSchemaVersion >= 64); Assert.Equal((long)DuckDbInitializer.CurrentSchemaVersion, Convert.ToInt64(await ScalarAsync(conn, "SELECT MAX(version) FROM schema_version"))); Assert.Equal(1L, Convert.ToInt64(await ScalarAsync(conn, "SELECT COUNT(*) FROM duckdb_indexes() WHERE table_name = 'query_store_health'"))); diff --git a/Lite/Database/DuckDbInitializer.cs b/Lite/Database/DuckDbInitializer.cs index 58b06614b..68a88b871 100644 --- a/Lite/Database/DuckDbInitializer.cs +++ b/Lite/Database/DuckDbInitializer.cs @@ -349,7 +349,7 @@ public void Dispose() /// /// Current schema version. Increment this when schema changes require table rebuilds. /// - internal const int CurrentSchemaVersion = 64; + internal const int CurrentSchemaVersion = 65; private readonly string _archivePath; @@ -2310,6 +2310,50 @@ await ExecuteNonQueryAsync(connection, } } } + + if (fromVersion < 65) + { + /* v65 (#4475, twinning Darling's V151): ag_replica_states and ag_database_replica_states gain + group_id — the GUID sys.availability_groups.group_id the engine stamps identically on every + replica of one Availability Group, stored as text (the collector column vocabulary has no uuid + type; AgDatabaseReplicaStatesCollector's last_hardened_lsn/last_commit_lsn already store a wide + identifier the same way). AgTopology.CountDistinctGroups uses it to close the one gap the + name-plus-replica-overlap rule (#4475) could not: two monitored SECONDARIES of one AG, with its + primary unmonitored, share no replica name with each other and so counted as two groups. A row + carrying group_id groups by it exactly; a row from before this rung carries none and falls back + to the pre-#4475 name + overlap rule. + + Appended at the end of each PayloadColumns list, so the positional appender and old parquet are + unaffected. Nothing to backfill and nothing that COULD be: a row collected before the upgrade + never asked the engine for its AG's group_id, and NULL is the honest value — a reader treats it + as "fall back to the pre-#4475 rule", exactly today's behavior. + + REQUIRED on this side for the v60 reason: the appender writes one value per declared payload + column, so a database without the column fails EndRow() on the first AG-collector batch — the + whole batch, not the column. Fresh installs get both from DuckDbSchemaGenerator (the AG tables + are generated from the shared collector catalog, not hand-written here); these ALTERs are for + an existing database and are idempotent. Neither AG table has ever had a v_ passthrough view + (view-less since Darling's V34 twin), so this rung is two ALTERs and nothing else. Non-fatal + per statement, matching v59–v64. */ + _logger?.LogInformation("Running migration to v65: ag_replica_states and ag_database_replica_states store each Availability Group's engine-assigned id, closing a replica-name-overlap gap in the distinct-group count"); + + foreach (var table in new[] + { + "ag_replica_states", + "ag_database_replica_states", + }) + { + try + { + await ExecuteNonQueryAsync(connection, + $"ALTER TABLE {table} ADD COLUMN IF NOT EXISTS group_id VARCHAR"); + } + catch (Exception ex) + { + _logger?.LogWarning("Migration to v65 on {Table}.group_id encountered an error (non-fatal): {Error}", table, ex.Message); + } + } + } } /// diff --git a/Lite/Services/LocalDataService.AgTopology.cs b/Lite/Services/LocalDataService.AgTopology.cs index 581605281..c2f374175 100644 --- a/Lite/Services/LocalDataService.AgTopology.cs +++ b/Lite/Services/LocalDataService.AgTopology.cs @@ -50,7 +50,8 @@ public partial class LocalDataService synchronization_health_desc, availability_mode_desc, failover_mode_desc, - endpoint_url + endpoint_url, + group_id FROM ag_replica_states AS r WHERE r.collection_time = (SELECT MAX(x.collection_time) FROM ag_replica_states AS x WHERE x.server_id = r.server_id) ORDER BY r.server_name, r.ag_name, r.replica_server_name"; @@ -75,7 +76,8 @@ FROM ag_replica_states AS r is_suspended, suspend_reason_desc, availability_mode_desc, - secondary_lag_seconds + secondary_lag_seconds, + group_id FROM ag_database_replica_states AS d WHERE d.collection_time = (SELECT MAX(x.collection_time) FROM ag_database_replica_states AS x WHERE x.server_id = d.server_id) ORDER BY d.server_name, d.ag_name, d.database_name, d.replica_server_name"; @@ -122,6 +124,7 @@ private async Task> ReadAgTopologyReplicasAsync() AvailabilityModeDesc = AgText(reader, 11), FailoverModeDesc = AgText(reader, 12), EndpointUrl = AgText(reader, 13), + GroupId = AgText(reader, 14), }); } @@ -157,6 +160,7 @@ private async Task> ReadAgTopologyDatabasesAsync() SuspendReasonDesc = AgText(reader, 13), AvailabilityModeDesc = AgText(reader, 14), SecondaryLagSeconds = AgCount(reader, 15), + GroupId = AgText(reader, 16), }); } diff --git a/PerformanceMonitor.Collectors/AgDatabaseReplicaStatesCollector.cs b/PerformanceMonitor.Collectors/AgDatabaseReplicaStatesCollector.cs index debbcc236..113b292fb 100644 --- a/PerformanceMonitor.Collectors/AgDatabaseReplicaStatesCollector.cs +++ b/PerformanceMonitor.Collectors/AgDatabaseReplicaStatesCollector.cs @@ -122,7 +122,8 @@ public readonly record struct Row( DateTime? LastRedoneTime, DateTime? LastReceivedTime, double? EstRedoCompletionTimeMin, - double? EstSendDrainTimeMin); + double? EstSendDrainTimeMin, + string? GroupId); private const string QueryText = @" SET TRANSACTION ISOLATION LEVEL READ UNCOMMITTED; @@ -148,7 +149,8 @@ public readonly record struct Row( last_redone_time = hdrs.last_redone_time, last_received_time = hdrs.last_received_time, est_redo_completion_time_min = CONVERT(float, (hdrs.redo_queue_size * 1.0 / NULLIF(hdrs.redo_rate, 0)) / 60.0), - est_send_drain_time_min = CONVERT(float, (hdrs.log_send_queue_size * 1.0 / NULLIF(hdrs.log_send_rate, 0)) / 60.0) + est_send_drain_time_min = CONVERT(float, (hdrs.log_send_queue_size * 1.0 / NULLIF(hdrs.log_send_rate, 0)) / 60.0), + group_id = CONVERT(nvarchar(36), ag.group_id) FROM sys.dm_hadr_database_replica_states AS hdrs JOIN sys.availability_replicas AS ar ON hdrs.replica_id = ar.replica_id @@ -206,6 +208,9 @@ are last here too. */ new CollectorColumn("last_received_time", CollectorColumnType.Timestamp), new CollectorColumn("est_redo_completion_time_min", CollectorColumnType.Double), new CollectorColumn("est_send_drain_time_min", CollectorColumnType.Double), + /* Appended LAST (#4475's group-count fallback rung), same shape and reasoning as the replica-grain + twin in AgReplicaStatesCollector: sys.availability_groups.group_id, stored as text. */ + new CollectorColumn("group_id", CollectorColumnType.Varchar), }; public override async ValueTask> ReadAsync(DbDataReader reader, CollectorContext context, CancellationToken cancellationToken) @@ -235,7 +240,8 @@ public override async ValueTask> ReadAsync(DbDataReader reader, Collec LastRedoneTime: reader.IsDBNull(17) ? null : reader.GetDateTime(17), LastReceivedTime: reader.IsDBNull(18) ? null : reader.GetDateTime(18), EstRedoCompletionTimeMin: reader.IsDBNull(19) ? null : reader.GetDouble(19), - EstSendDrainTimeMin: reader.IsDBNull(20) ? null : reader.GetDouble(20))); + EstSendDrainTimeMin: reader.IsDBNull(20) ? null : reader.GetDouble(20), + GroupId: reader.IsDBNull(21) ? null : reader.GetString(21))); } return rows; @@ -264,6 +270,7 @@ public override void WritePayload(Row row, ICollectorRowWriter writer, Collector .Value(row.LastRedoneTime) /* last_redone_time TIMESTAMP */ .Value(row.LastReceivedTime) /* last_received_time TIMESTAMP */ .Value(row.EstRedoCompletionTimeMin) /* est_redo_completion_time_min DOUBLE (min) */ - .Value(row.EstSendDrainTimeMin); /* est_send_drain_time_min DOUBLE (min) */ + .Value(row.EstSendDrainTimeMin) /* est_send_drain_time_min DOUBLE (min) */ + .Value(row.GroupId); /* group_id VARCHAR (uuid text, #4475) */ } } diff --git a/PerformanceMonitor.Collectors/AgReplicaStatesCollector.cs b/PerformanceMonitor.Collectors/AgReplicaStatesCollector.cs index c75bbc4b1..fefa7d337 100644 --- a/PerformanceMonitor.Collectors/AgReplicaStatesCollector.cs +++ b/PerformanceMonitor.Collectors/AgReplicaStatesCollector.cs @@ -72,7 +72,8 @@ public readonly record struct Row( string? AvailabilityModeDesc, string? FailoverModeDesc, string? EndpointUrl, - bool? IsLocal); + bool? IsLocal, + string? GroupId); private const string QueryText = @" SET TRANSACTION ISOLATION LEVEL READ UNCOMMITTED; @@ -88,7 +89,8 @@ public readonly record struct Row( availability_mode_desc = ar.availability_mode_desc, failover_mode_desc = ar.failover_mode_desc, endpoint_url = ar.endpoint_url, - is_local = ars.is_local + is_local = ars.is_local, + group_id = CONVERT(nvarchar(36), ag.group_id) FROM sys.availability_replicas AS ar JOIN sys.availability_groups AS ag ON ar.group_id = ag.group_id @@ -131,6 +133,12 @@ ORDER BY new CollectorColumn("failover_mode_desc", CollectorColumnType.Varchar), new CollectorColumn("endpoint_url", CollectorColumnType.Varchar), new CollectorColumn("is_local", CollectorColumnType.Boolean), + /* Appended LAST (#4475's group-count fallback rung): sys.availability_groups.group_id, the SAME + GUID on every replica of one AG, stored as text (36 chars) rather than a uniqueidentifier — the + collector vocabulary has no uuid/GUID type, and every consumer treats it as an opaque identity + string rather than doing arithmetic on it, same reasoning as last_hardened_lsn/last_commit_lsn on + the sibling collector. */ + new CollectorColumn("group_id", CollectorColumnType.Varchar), }; public override async ValueTask> ReadAsync(DbDataReader reader, CollectorContext context, CancellationToken cancellationToken) @@ -150,7 +158,8 @@ public override async ValueTask> ReadAsync(DbDataReader reader, Collec AvailabilityModeDesc: reader.IsDBNull(7) ? null : reader.GetString(7), FailoverModeDesc: reader.IsDBNull(8) ? null : reader.GetString(8), EndpointUrl: reader.IsDBNull(9) ? null : reader.GetString(9), - IsLocal: reader.IsDBNull(10) ? null : reader.GetBoolean(10))); + IsLocal: reader.IsDBNull(10) ? null : reader.GetBoolean(10), + GroupId: reader.IsDBNull(11) ? null : reader.GetString(11))); } return rows; @@ -169,6 +178,7 @@ public override void WritePayload(Row row, ICollectorRowWriter writer, Collector .Value(row.AvailabilityModeDesc) /* availability_mode_desc VARCHAR */ .Value(row.FailoverModeDesc) /* failover_mode_desc VARCHAR */ .Value(row.EndpointUrl) /* endpoint_url VARCHAR */ - .Value(row.IsLocal); /* is_local BOOLEAN */ + .Value(row.IsLocal) /* is_local BOOLEAN */ + .Value(row.GroupId); /* group_id VARCHAR (uuid text, #4475) */ } } diff --git a/PerformanceMonitor.Common/AgTopology.cs b/PerformanceMonitor.Common/AgTopology.cs index a0bd2598c..2e3076d23 100644 --- a/PerformanceMonitor.Common/AgTopology.cs +++ b/PerformanceMonitor.Common/AgTopology.cs @@ -211,6 +211,10 @@ public static List BuildCards( DatabaseCollectionTime = dbRows is { Count: > 0 } ? dbRows[0].CollectionTime : null, PrimaryReplica = replicaItems.FirstOrDefault(r => r.IsPrimary)?.ReplicaServerName, Severity = worst, + /* Every replica row of one AG carries the SAME group_id (the engine stamps one id per AG), so + the first non-null one found is the card's — tolerant of a row or two still NULL mid-upgrade + rather than requiring every row to agree. */ + GroupId = replicaItems.Select(r => r.GroupId).FirstOrDefault(g => !string.IsNullOrWhiteSpace(g)), }; /* Replicas/Databases are get-only ObservableCollections (#4238) — populated here, once, rather than @@ -272,30 +276,59 @@ public static (int DistinctGroups, int ReportingServers, int Views) Counts(IRead ArgumentNullException.ThrowIfNull(cards); return ( - CountDistinctGroups(cards.Select(c => (c.AgName, (IEnumerable)c.Replicas.Select(r => r.ReplicaServerName)))), + CountDistinctGroups(cards.Select(c => (c.AgName, (IEnumerable)c.Replicas.Select(r => r.ReplicaServerName), c.GroupId))), cards.Select(c => c.ServerId).Distinct().Count(), cards.Count); } /// - /// Counts distinct AG groups by connected components (#4475): two members with the same name - /// (case-insensitive) union into one group when their replica-name sets overlap (share at least one name, - /// case-insensitive); a member with an EMPTY replica set unions ONLY with other same-named EMPTY members - /// (name-only matching among themselves), never with a same-named member that has replicas, since nothing - /// on an empty row ties it to one specific AG over another. Members with - /// different names never union, regardless of their replica sets. Shared by and the - /// MCP/web AG reader's equivalent distinct-AG count (DarlingAgReader.Build), so the two surfaces - /// cannot drift back apart. No store read: both callers already carry the replica names on the rows they - /// group. + /// Counts distinct AG groups by connected components (#4475), with a members overload that carries no + /// group_id — every member falls back to the name-plus-overlap rule below, unchanged from before + /// V151. Kept so the pre-existing id-less callers and pins compile and behave exactly as before. /// public static int CountDistinctGroups(IEnumerable<(string? AgName, IEnumerable ReplicaServerNames)> members) { ArgumentNullException.ThrowIfNull(members); + return CountDistinctGroups(members.Select(m => (m.AgName, m.ReplicaServerNames, (string?)null))); + } + + /// + /// Counts distinct AG groups by connected components, extended for sys.availability_groups.group_id + /// (#4475/V151). The count rule: + /// + /// A member that carries a -shaped GroupId (case-insensitive text compare) + /// unions with every OTHER member whose GroupId matches it EXACTLY — the engine stamps the same GUID + /// on every replica of one AG, so this closes the one gap the name-plus-overlap rule alone could not: two + /// monitored SECONDARIES of one AG, with its primary unmonitored, share no replica name with each other + /// (each reports only itself under sys.dm_hadr_availability_replica_states's local-only rule), so + /// they union on group_id alone even though their replica-name sets are disjoint. + /// Two members that BOTH carry a group_id, but DIFFERENT ones, never union — even when they share a + /// name and their replica sets overlap. A group_id is definitive: two different ids can never mean the same + /// AG, so nothing below overrides that verdict. + /// A member with NO group_id (a row collected before V151, or an id-less caller) falls back to the + /// pre-existing name-plus-replica-overlap rule (below) among the OTHER id-less members — case-insensitive + /// name match, replica-name sets overlap (or both empty, matched by name alone), never bridging an empty + /// set to a same-named non-empty one. + /// A WITH-id member and a WITHOUT-id member of the SAME name whose replica sets overlap still union — + /// the same AG, observed once before the V151 upgrade landed on that reporter and once after. The id-less + /// member does not get a group_id from this union; it simply joins the same connected component. + /// + /// Shared by and the MCP/web AG reader's equivalent distinct-AG count + /// (DarlingAgReader.Build), so the two surfaces cannot drift back apart. No store read: both callers + /// already carry the replica names (and now the group id) on the rows they group. + /// + public static int CountDistinctGroups(IEnumerable<(string? AgName, IEnumerable ReplicaServerNames, string? GroupId)> members) + { + ArgumentNullException.ThrowIfNull(members); + var items = members - .Select(m => (Name: Key(m.AgName), Replicas: new HashSet( - m.ReplicaServerNames.Select(n => (n ?? "").ToUpperInvariant()).Where(n => n.Length > 0), - StringComparer.Ordinal))) + .Select(m => ( + Name: Key(m.AgName), + Replicas: new HashSet( + m.ReplicaServerNames.Select(n => (n ?? "").ToUpperInvariant()).Where(n => n.Length > 0), + StringComparer.Ordinal), + GroupId: string.IsNullOrWhiteSpace(m.GroupId) ? null : m.GroupId.Trim().ToUpperInvariant())) .ToList(); var parent = new int[items.Count]; @@ -324,8 +357,39 @@ void Union(int a, int b) } } - /* Only same-named members can ever union, so grouping by name first keeps the pairwise comparison - quadratic within one AG name's cards rather than across the whole fleet. */ + /* Pass 1: members that carry a group_id union EXACTLY on it, regardless of name -- a group_id is + definitive, and two different ids never union even under a matching name (handled by never reaching + the name+overlap pass below for a with-id member paired with another with-id member of a DIFFERENT + id). */ + var byGroupId = new Dictionary>(StringComparer.Ordinal); + for (var i = 0; i < items.Count; i++) + { + if (items[i].GroupId is string groupId) + { + if (!byGroupId.TryGetValue(groupId, out var indices)) + { + indices = new List(); + byGroupId[groupId] = indices; + } + + indices.Add(i); + } + } + + foreach (var indices in byGroupId.Values) + { + for (var a = 1; a < indices.Count; a++) + { + Union(indices[0], indices[a]); + } + } + + /* Pass 2: the pre-existing name-plus-overlap rule, run over every pair that is NOT both with-id-and- + different-id. A with-id member still takes part here so it can join a same-named WITHOUT-id member + whose replicas overlap (the same AG seen before and after the V151 upgrade); two with-id members of + DIFFERENT ids must never union even if their names match and replicas overlap, so that specific pair + is skipped. Only same-named members can ever union, so grouping by name first keeps the pairwise + comparison quadratic within one AG name's cards rather than across the whole fleet. */ var byName = new Dictionary>(StringComparer.Ordinal); for (var i = 0; i < items.Count; i++) { @@ -345,6 +409,16 @@ quadratic within one AG name's cards rather than across the whole fleet. */ for (var b = a + 1; b < indices.Count; b++) { var (ia, ib) = (indices[a], indices[b]); + + /* Two members that BOTH carry a group_id never union here on name+overlap alone -- pass 1 + already decided their relationship definitively (same id -> already unioned; different + id -> must never union), and letting name+overlap override that would defeat the whole + point of carrying an id. */ + if (items[ia].GroupId is not null && items[ib].GroupId is not null) + { + continue; + } + var setA = items[ia].Replicas; var setB = items[ib].Replicas; @@ -405,6 +479,7 @@ private static AgTopologyReplica ToReplica(AgTopologyReplicaRow row) FailoverModeDesc = row.FailoverModeDesc, EndpointUrl = row.EndpointUrl, Severity = Worse(Worse(Worse(Worse(syncHealth, connected), operational), recovery), role), + GroupId = row.GroupId, }; } @@ -572,6 +647,10 @@ public sealed class AgTopologyReplicaRow public string? AvailabilityModeDesc { get; init; } public string? FailoverModeDesc { get; init; } public string? EndpointUrl { get; init; } + + /// sys.availability_groups.group_id as text (V151, #4475) — the same GUID on every replica + /// of one AG, stamped by the engine. Null on a row collected before this column existed. + public string? GroupId { get; init; } } /// One database-grain row as collected. Queue sizes are KB and rates KB/s (the DMV's units), both @@ -594,6 +673,11 @@ public sealed class AgTopologyDatabaseRow public string? SuspendReasonDesc { get; init; } public string? AvailabilityModeDesc { get; init; } public long? SecondaryLagSeconds { get; init; } + + /// sys.availability_groups.group_id as text (V151, #4475), carried on the database grain + /// too since the collector stamps it there as well. Unused by today (the + /// count rule works off the replica grain), kept for symmetry with the SQL and the storage-layer row. + public string? GroupId { get; init; } } /// @@ -631,6 +715,10 @@ public sealed class AgTopologyReplica : AgTopologyObservable public string? EndpointUrl { get; set; } public HealthSeverity Severity { get; set; } + /// sys.availability_groups.group_id as text (V151, #4475). Null on a row collected before + /// this column existed — the count rule falls back to name+overlap for those. + public string? GroupId { get; set; } + public string RoleDisplay => string.IsNullOrWhiteSpace(RoleDesc) ? "UNKNOWN ROLE" : RoleDesc!; /// @@ -695,6 +783,7 @@ public void UpdateFrom(AgTopologyReplica latest) FailoverModeDesc = latest.FailoverModeDesc; EndpointUrl = latest.EndpointUrl; Severity = latest.Severity; + GroupId = latest.GroupId; RaiseAllPropertiesChanged(); } @@ -784,6 +873,12 @@ public sealed class AgTopologyCard : AgTopologyObservable public string? PrimaryReplica { get; set; } public HealthSeverity Severity { get; set; } + /// sys.availability_groups.group_id as text (V151, #4475) — taken from the card's first + /// replica row, which is the same value every replica of this AG stores (the engine stamps one group_id + /// per AG, identically on every replica). Null on a card built entirely from rows collected before this + /// column existed. + public string? GroupId { get; set; } + /// Get-only and never reassigned (#4238): the SAME collection instance lives for the card's whole /// life, so the nested, non-virtualized replica-chip ItemsControl never sees a new /// ItemsSource reference and never rebuilds its containers either. mutates @@ -844,6 +939,7 @@ public void UpdateFrom(AgTopologyCard latest) DatabaseCollectionTime = latest.DatabaseCollectionTime; PrimaryReplica = latest.PrimaryReplica; Severity = latest.Severity; + GroupId = latest.GroupId; AgTopology.Reconcile( Replicas, latest.Replicas,